davidzollo opened a new pull request, #12110:
URL: https://github.com/apache/seatunnel/pull/12110

   ## Purpose
   
   This E2E test documents Zeta's actual, current behavior when 
`--restore-with-checkpoint` is asked to restore from a job whose **latest** 
on-disk checkpoint file is corrupted, while an **older** checkpoint (still 
within `checkpoint.storage.max-retained`) remains intact and untouched.
   
   Verified by reading current `dev` HEAD source (not assumed):
   
   - Checkpoint restore reads the latest checkpoint via the configured storage 
plugin (`LocalFileStorage` by default) and deserializes it with **no 
checksum/CRC validation anywhere in the read path**.
   - `LocalFileStorage.getCheckpointsByJobIdAndPipelineId` only catches 
`IOException` per file
     
(`checkpoint-storage-plugins/checkpoint-storage-local-file/.../LocalFileStorage.java`),
 but a corrupted file's deserialization failure comes out of 
`ProtoStuffSerializer.deserialize` as an **unchecked** exception (its method 
signature carries no `throws` clause at all — verified both by inspecting 
`serializer-protobuf/.../ProtoStuffSerializer.java` and by an isolated, 
non-Docker reproduction described below), so the failure is never swallowed.
   - The exception propagates out of `CheckpointManager`'s constructor (via 
`ExceptionUtil.sneakyThrow` in `CheckpointManager.java`), out of 
`JobMaster.initCheckPointManager()`, and is caught by `JobMaster.init()`.
   - For a brand-new job submission via `--restore-with-checkpoint` (as opposed 
to a master-failover re-init), `JobMaster.init(..., restart=false)` is what 
`CoordinatorService.submitJob` calls (`CoordinatorService.java`), so 
`cancelJob()` itself is not invoked — but `init()` unconditionally rethrows the 
captured exception, which fails `submitJob`'s `jobSubmitFuture`, which fails 
the CLI's job submission with a non-zero exit code.
   - There is **no fallback to the older, intact checkpoint**, even though 
`checkpoint.storage.max-retained: 3` (this module's test config) keeps several 
around.
   
   This is Tier 2 of the ongoing Zeta-engine-core E2E initiative: a 
verified-fragile behavior of the current architecture with no existing 
regression coverage, documented as today's baseline (not a claim that this is 
the *ideal* behavior — a future fix adding checksum validation + fallback would 
rightly flip this test's outcome, and this test would need updating alongside 
that fix).
   
   ## What the test does
   
   
`CheckpointRestoreWithStopIT#testRestoreFailsWhenLatestCheckpointFileIsCorrupted`:
   
   1. Runs a real streaming job (`CheckpointableSequenceSource` -> `LocalFile` 
sink) to **2 completed checkpoints** (`checkpoint.interval = 3000`, 
`checkpoint.retain-after-job-cancelled = true`), then stops it cleanly so both 
checkpoint files survive on disk.
   2. Lists the job's checkpoint files directly inside the container 
(`/tmp/seatunnel/checkpoint_snapshot/<jobId>/`, matching this module's 
configured `namespace`) and replicates `AbstractCheckpointStorage`'s own 
"latest" selection rule (largest leading `<epochMillis>` segment of the 
`<epochMillis>-<random>-<pipelineId>-<checkpointId>.ser` file name) to find 
precisely the file production code would read back.
   3. Corrupts **only** that latest file, leaving the older one byte-for-byte 
untouched (checksummed before/after to prove it).
   4. Attempts `--restore-with-checkpoint` against the original (now-terminal) 
job id, bounded by a 2-minute timeout so a hang would fail the test loudly 
instead of blocking forever.
   5. Asserts the restore attempt fails with a non-zero exit code and a 
non-empty diagnostic, that no new sink rows appear (no silent 
restart-from-scratch and no silent fallback), and that the untouched older 
checkpoint file(s) remain byte-for-byte identical.
   
   ### How the corruption is constructed, and why
   
   The corruption overwrites the file in place with random bytes of the **same 
length** (no truncation — this documents same-size on-disk corruption, e.g. a 
bad block or a crash mid-write, not a truncated file), then forces the file's 
**first byte** to `0x00`.
   
   That last part matters and is not incidental. Protostuff/protobuf's wire 
format reads a leading varint field tag starting at byte 0, and field number 
`0` is a reserved, always-invalid tag every decoder rejects immediately — 
before any other byte is even interpreted. Plain full-file random overwrite 
alone does **not** reliably reproduce this test's documented failure: I built a 
standalone, non-Docker reproduction (`RuntimeSchema`/`ProtostuffIOUtil` against 
a POJO with the same field shape as `PipelineState` — `String, int, long, 
byte[]`) and found random bytes form a self-consistent, non-throwing protostuff 
field sequence on roughly 1 in 5 attempts, independent of file size (sampled 
sizes 8 through 2014 bytes, 200 trials each, ~78-82% throw rate — never 100%). 
That would make a pure-random-bytes version of this test flaky. Forcing only 
the first byte to an invalid tag while leaving every other byte genuinely 
random keeps the corruption realistic (e.g. a torn write or block-ze
 roing landing on the start of the file) while making the deserialization 
failure deterministic: the same reproduction threw on 3600/3600 trials across 
file sizes from 1 to 2014 bytes once the leading byte was forced to `0x00`.
   
   ## Test plan
   
   - `./mvnw spotless:apply -pl 
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -am -nsu` — 
BUILD SUCCESS.
   - `./mvnw install -pl 
'seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base,!seatunnel-engine/seatunnel-engine-ui'
 -am -nsu -Dmaven.gitcommitid.skip=true -DskipTests -Dspotless.check.skip=true 
-T 3C` — BUILD SUCCESS; confirmed the compiled test class exists under 
`target/test-classes`.
   - Per this repository's local-verification policy for Apache SeaTunnel, the 
new E2E test itself (Docker/Testcontainers-based) is not executed locally in 
this change; it runs as part of this module's existing E2E suite in CI. The 
behavioral claims above are substantiated by direct inspection of `dev` HEAD 
source across the full call chain (`ClientCommandArgs` -> `SeaTunnelClient` -> 
`CheckpointManager` -> `LocalFileStorage`/`AbstractCheckpointStorage` -> 
`ProtoStuffSerializer` -> `JobMaster.init()` -> `CoordinatorService.submitJob`) 
plus the standalone protostuff reproduction described above.
   - No changes to `src/main/**`; this is a test-only addition to 
`CheckpointRestoreWithStopIT.java`.
   


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