feat: DAL-99 add transactional sink capabilities - #111
Conversation
Up to standards ✅🟢 Issues
|
| Metric | Results |
|---|---|
| Complexity | 25 |
| Duplication | 2 |
NEW Get contextual insights on your PRs based on Codacy's metrics, along with PR and Jira context, without leaving GitHub. Enable AI reviewer
TIP This summary will be updated as you push new changes.
There was a problem hiding this comment.
Pull request overview
This PR extends the streaming runtime’s sink model to represent ordinary vs transactional vs epoch-idempotent behavior at binding construction time, introduces stable sink IDs, and tightens checkpoint/recovery handling to surface durable-commit uncertainty as RecoveryRequired while redacting pre-commit metadata from public diagnostics.
Changes:
- Introduce
StableSinkIdand enforce uniqueness (including rejecting reuse across outputs) during whole-job preflight. - Replace dynamic sink delivery probing with a structurally fixed
SinkCapabilityonOrdinarySinkBinding, and propagate that into manifests/metrics. - Update checkpoint commit/recovery paths and soak tests to treat partial sink commit uncertainty as
RecoveryRequired, plus add diagnostic redaction tests.
Reviewed changes
Copilot reviewed 4 out of 4 changed files in this pull request and generated 1 comment.
| File | Description |
|---|---|
| crates/calc-flow/src/runtime/streaming/soak.rs | Adjust soak matrix expectations for partial sink commit faults to assert RecoveryRequired. |
| crates/calc-flow/src/runtime/streaming/sink_task.rs | Add pipeline identity to sink tasks, redact abort/recovery diagnostics, and emit RecoveryRequired for durable-commit uncertainty. |
| crates/calc-flow/src/runtime/streaming/runner.rs | Thread pipeline identity into sink tasks, adopt stable sink IDs in connector resource tracking, and classify sink checkpoint failures for job state. |
| crates/calc-flow/src/runtime/streaming/job.rs | Add StableSinkId and SinkCapability, freeze capability at binding construction, and validate sink ID uniqueness across outputs. |
Suppressed comments (1)
crates/calc-flow/src/runtime/streaming/sink_task.rs:1067
recover_transactional_sinksonly treatsErrresults asRecoveryRequired; if a sinkrecover()panics, the panic will unwind past this function and can crash the runner (and potentially surface unredacted manifest/pre-commit data via panic payload/backtrace). Wrap therecover()call incatch_unwindand map bothErrand panics to the same redactedRecoveryRequiredoutcome.
if sink.binding.recover(manifest).await.is_err() {
return Err(CalcFlowError::RecoveryRequired {
pipeline_name: manifest.pipeline_name().into(),
message: format!(
"sink {:?} did not recover durable epoch {}; retry recovery before allocating a new epoch",
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
Make sink identity and capability structural before connector lifecycle work. Preserve durable post-manifest transactions for forward recovery and keep pre-commit metadata out of diagnostics. Co-authored-by: multica-agent <github@multica.ai>
ae4e9c4 to
c7470be
Compare
Use capability-neutral missing-sink errors and convert connector recovery panics into redacted forward-recovery failures. Co-authored-by: multica-agent <github@multica.ai>
wegamekinglc
left a comment
There was a problem hiding this comment.
Blocking review (recorded as COMMENTED because the authenticated GitHub identity is also the PR author and cannot submit REQUEST_CHANGES).
crates/calc-flow/src/runtime/streaming/sink_task.rs:1063: recovery panic payloads can still leak pre-commit metadata through the Rust panic hook. FutureExt::catch_unwind only changes the returned error after unwinding begins; the active panic hook runs first and normally writes the payload to stderr/logs. The new test sink demonstrates the exposure at lines 1307-1311 by panicking with PRECOMMIT_SENTINEL, but the assertion at lines 2218-2223 checks only the formatted CalcFlowError, so it cannot detect hook output. This violates DAL-99 acceptance that pre-commit metadata must not enter logs/public diagnostics.
Please make the recovery boundary guarantee that manifest/pre-commit data cannot reach panic-hook/log output and add a regression that observes the actual hook/log channel, not only the returned error. Avoid unsynchronized temporary replacement of the process-global panic hook, which would be racy with concurrent runtime tasks/tests. Re-review the resulting exact head before merge.
Summary
RecoveryRequiredwithout leaking pre-commit metadataCloses DAL-99
Closes #97
Test plan
uv run python scripts/run_rust_tests.pycargo fmt --all --checkcargo clippy --workspace --all-targets --all-features -- -D warningsRUSTDOCFLAGS="-D warnings" cargo doc --workspace --all-features --no-depscargo llvm-cov --workspace --all-features --fail-under-lines 90(90.40% lines)