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]
