SEZ9 opened a new issue, #11980:
URL: https://github.com/apache/seatunnel/issues/11980

   ### Search before asking
   
   - [X] I had searched in the 
[issues](https://github.com/apache/seatunnel/issues?q=is%3Aissue) and found no 
similar feature requirement.
   
   ### Description
   
   #### Motivation
   
   The REST API exposes a job's *current* status only. Three questions that 
come up
   first during incident triage cannot be answered from the engine:
   
   1. "The job shows RUNNING — has it silently restarted?"
   2. "Is a pipeline crash-looping?"
   3. "How long did it wait for resources before running?"
   
   Today all three require reading node logs. Nothing in the REST API or the
   Prometheus exporter answers them: `JobMetricExports` only publishes a 
`job_count`
   gauge by status — there is no restart or failover counter anywhere.
   
   Note on the push side: `PhysicalPlan#reportJobStateEvent` emits 
`JobStateEvent`
   only for **terminal** states unless `report-non-terminal-job-state` is 
enabled
   (`EngineConfig` default `false`), and `JobStateEvent` carries no 
`fromState`. So
   event push does not answer these questions in a default deployment either.
   
   #### What already exists (and is invisible)
   
   Two signals are already tracked and simply never surfaced:
   
   - **Pipeline restart count.** `SubPlan#prepareRestorePipeline()` increments
     `pipelineRestoreNum` on every pipeline restore and gates further restores 
on
     `pipelineMaxRestoreNum`. There is a getter (`getPipelineRestoreNum()`), 
used
     only for that gate and a log line.
   - **Per-state entry timestamps.** `IMAP_STATE_TIMESTAMPS` holds a `Long[]`
     indexed by state ordinal, for **both** job level (key: `jobId`) and 
pipeline
     level (key: `PipelineLocation`). REST currently reads exactly one slot of 
it —
     `SCHEDULED` — to render `startTime` (`BaseService#getJobStartTime`).
   
   Important: a pipeline restart does **not** change the job status. The job 
stays
   `RUNNING` across `SubPlan#reset()`. And at job level a crash-loop of the form
   `RUNNING → FAILING → RUNNING` is impossible: `FAILING` can only proceed to 
the
   terminal `FAILED`, and `PhysicalPlan#updateJobState` rejects leaving a 
terminal
   state. **Therefore job-level state history alone cannot answer question 2** —
   pipeline-level data is required. This is why the proposal is split in two 
phases.
   
   #### Phase 1 — surface the signals that already exist (no hot-path change)
   
   Add to `GET /job-info/{jobId}` a per-pipeline diagnostic block and the full 
set
   of already-recorded state timestamps:
   
   ```json
   {
     "jobId": "...",
     "jobStatus": "RUNNING",
     "diagnostics": {
       "stateTimestamps": {
         "CREATED": 1755000000000,
         "PENDING": 1755000000500,
         "SCHEDULED": 1755000001000,
         "RUNNING": 1755000003000
       },
       "pipelines": [
         { "pipelineId": 1, "status": "RUNNING", "restoreCount": 7, 
"maxRestoreCount": 100 }
       ]
     }
   }
   ```
   
   Properties of Phase 1:
   
   - **Pure read.** No new recording, no write on the state-transition path.
   - Restart count is read from the existing `SubPlan` counter; timestamps from 
the
     existing IMap, which any node can already read directly — no new Hazelcast
     operation needed.
   - Immediately answers questions 2 and 3, and makes "restarted since submit"
     visible for question 1.
   
   #### Phase 2 — bounded transition history (needs design agreement)
   
   Record `(fromState, toState, timestampMs)` into a bounded, latest-first 
buffer,
   at **both** job and pipeline level.
   
   - **Where it lives:** an `AuxiliaryStateStores` member — that bundle is
     documented as state "closer to observability, recent history, or cleanup 
than
     to failover correctness", which is exactly this data. Keying follows the
     existing timestamps map (`jobId` / `PipelineLocation`).
   - **Why a distributed store rather than JobMaster heap:**
     1. The transition path **already** performs one IMap write per transition
        (`updateStateTimestamps` → `runningJobStateTimestampsIMap.set(...)`), 
so a
        store-backed history introduces no *new class* of cost on that path.
     2. REST can be served by a non-master node; job-level timestamps are 
already
        read straight from the IMap by any node. A heap-only buffer would 
require a
        new operation just to read it back.
     3. A heap-only window is empty right after a master switch — losing 
exactly the
        history that motivates question 1.
   - **Capacity:** default 32, latest-first, oldest evicted — matching the
     checkpoint history buffer's behaviour 
(`PipelineCheckpointOverview#addHistory`).
     Note that 32 is hard-coded there
     (`new CheckpointMonitorService(engineContext, 32)`), not configurable; this
     proposal adds a config option so the feature can be disabled (`0`) and 
capped
     at a sane maximum.
   - **Cumulative counters** that survive eviction: `totalTransitions`,
     `totalPipelineRestores`, `firstCreatedTimestamp`.
   
   #### Safety design
   
   1. **Bounded growth.** Fixed capacity per key; a crash-looping job stays 
under
      ~1 KB. Capacity validated at startup.
   2. **Failure isolation.** The append is wrapped so any exception is caught,
      logged at WARN (rate-limited) and swallowed. A broken history degrades to
      "history unavailable", never to "transition failed".
   3. **Hot-path cost.** The append happens inside the already-`synchronized`
      `PhysicalPlan#updateJobState` / `SubPlan#updatePipelineState` — no new 
lock.
      It must be benchmarked against the existing per-transition IMap write; if
      measurable, coalescing history into the same write is the fallback.
   4. **Entry lifecycle.** The new store must be cleaned in the same paths that
      already do `removeKeys(runningJobStateTimestampsIMap, ...)`
      (`cleanupPendingJobStateForRestore`, terminal-zombie cleanup), otherwise 
every
      finished job leaks an entry. Explicit acceptance criterion.
   5. **Serialization compatibility.** A new store rather than widening the 
existing
      `Long[]` value type, so rolling upgrades are unaffected.
   
   #### Non-goals
   
   - Not an event store — long-term history stays with `event-report-http`.
   - No failover *cause* capture (exception context) — separate issue.
   - No new REST endpoint; no Web UI work here.
   
   #### Downstream consumers
   
   - Web UI timeline / "restarted N times" badge (follow-up).
   - CLI runtime diagnostics (#11616): turns log scraping into a small 
structured
     document.
   - The accuracy benchmark (#11553) currently judges a streaming job by "still
     alive after 60s". A pipeline that keeps restoring passes that check while 
the
     job status stays `RUNNING`. With `restoreCount` the criterion becomes 
"alive
     AND no restarts" — a concrete correctness fix for an existing consumer.
   
   #### Validation
   
   - Unit: eviction order and counter correctness under eviction; append-failure
     isolation (inject an exception, assert the transition still completes); 
cleanup
     removes entries for finished jobs.
   - Restart path: drive repeated pipeline restores (see
     `CheckpointErrorRestoreEndTest` for the existing pattern) and assert
     `restoreCount` reaches the expected value and is visible over REST.
   - E2E: field present in `/job-info/{jobId}`; disabled config removes the 
history
     block; `rest-api-v2.md` updated.
   
   #### Compatibility
   
   Purely additive response fields. Phase 1 changes no write path at all.
   
   ### Usage Scenario
   
   _No response_
   
   ### Related issues
   
   _No response_
   
   ### Are you willing to submit a PR?
   
   - [X] Yes I am willing to submit a PR!
   
   ### Code of Conduct
   
   - [X] I agree to follow this project's [Code of 
Conduct](https://www.apache.org/foundation/policies/conduct)
   


-- 
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