laskoviymishka commented on code in PR #3265:
URL: https://github.com/apache/iceberg-rust/pull/3265#discussion_r4153186875
##########
crates/iceberg/public-api.txt:
##########
@@ -3238,25 +3238,33 @@ pub struct iceberg::transaction::AddColumn
impl iceberg::transaction::AddColumn
pub fn iceberg::transaction::AddColumn::optional(name: impl
alloc::string::ToString, field_type: iceberg::spec::Type) -> Self
pub fn iceberg::transaction::AddColumn::required(name: impl
alloc::string::ToString, field_type: iceberg::spec::Type, initial_default:
iceberg::spec::Literal) -> Self
+impl core::clone::Clone for iceberg::transaction::AddColumn
Review Comment:
Deriving `Clone` on all eight actions is needed for `fresh_clone`, but it
also lands `Clone` in the public API for each — semver surface we can't easily
walk back, and it lets a caller clone a half-built action and `apply` both
copies (a `with_check_duplicate(false)` append cloned and applied twice
silently doubles data files). I'd at least doc-comment that these `Clone` impls
exist for transaction retry infrastructure, not user re-application.
##########
crates/iceberg/src/transaction/action.rs:
##########
@@ -25,28 +24,147 @@ use crate::table::Table;
use crate::transaction::Transaction;
use crate::{Result, TableRequirement, TableUpdate};
-/// A boxed, thread-safe reference to a `TransactionAction`.
-pub(crate) type BoxedTransactionAction = Arc<dyn TransactionAction>;
+/// A boxed entry pairing a transaction action with its retry-persistent state.
+pub(crate) type TransactionActionEntry = Box<dyn ErasedActionEntry>;
/// A trait representing an atomic action that can be part of a transaction.
///
/// Implementors of this trait define how a specific action is committed to a
table.
/// Each action is responsible for generating the updates and requirements
needed
/// to modify the table metadata.
+///
+/// An action's intent is immutable once applied to a transaction.
Retry-persistent
+/// state lives in the associated [`TransactionAction::State`], which survives
replay
+/// attempts of one logical execution but is never shared between executions
+/// (cloning a transaction creates fresh state via
[`TransactionAction::new_state`]).
#[async_trait]
-pub(crate) trait TransactionAction: AsAny + Sync + Send {
+pub(crate) trait TransactionAction: Clone + Send + Sync + 'static {
+ /// Retry-persistent state exclusively owned by one logical execution of
this
+ /// action. Stateless actions use `State = ()`.
+ type State: Send + Sync + 'static;
+
+ /// Creates fresh state for one logical execution.
+ ///
+ /// This is infallible and table-independent; table-dependent
initialization
+ /// happens during [`TransactionAction::commit`].
+ fn new_state(&self) -> Self::State;
+
/// Commits this action against the provided table and returns the
resulting updates.
/// NOTE: This function is intended for internal use only and should not
be called directly by users.
///
+ /// One replay attempt: the action (intent) is borrowed immutably, its
execution
+ /// state mutably, and the table reflects the current transaction-local
base.
+ ///
/// # Arguments
///
+ /// * `state` - The retry-persistent state for this logical execution.
/// * `table` - The current state of the table this action should apply to.
///
/// # Returns
///
/// An `ActionCommit` containing table updates and table requirements,
/// or an error if the commit fails.
- async fn commit(self: Arc<Self>, table: &Table) -> Result<ActionCommit>;
+ async fn commit(&self, state: &mut Self::State, table: &Table) ->
Result<ActionCommit>;
+
+ /// Best-effort cleanup after the transaction reaches a terminal result.
+ ///
+ /// Consumes the action and its state: the type system guarantees that no
+ /// further attempt can run for this entry after cleanup. Cleanup must not
+ /// change the already-determined transaction result.
+ ///
+ /// See [`CommitStatus`] for what each terminal status allows cleanup to
do.
+ // TODO: invoke this from `Transaction::commit` once terminal status
+ // classification is wired up (stateful transaction RFC, milestone 3).
+ #[allow(dead_code)]
+ async fn cleanup(self: Box<Self>, state: Self::State, table: &Table,
status: CommitStatus);
Review Comment:
`cleanup` being required means all eight actions carry an identical empty
body, and the compiler can't tell a deliberate no-op from a forgotten one. I'd
give it a default empty impl on the trait so stateless actions inherit it and
only stateful ones override.
That also lines up with the `StatelessAction` direction @Stefan-Dienst
floated — a blanket `TransactionAction` impl over a `StatelessAction: Clone +
Send + Sync + 'static` that only defines `commit` would collapse the
`State`/`new_state`/`cleanup` boilerplate for the eight stateless actions down
to a single method. Since the trait is `pub(crate)`, the blanket impl has no
public-API cost.
##########
crates/iceberg/src/transaction/action.rs:
##########
@@ -25,28 +24,147 @@ use crate::table::Table;
use crate::transaction::Transaction;
use crate::{Result, TableRequirement, TableUpdate};
-/// A boxed, thread-safe reference to a `TransactionAction`.
-pub(crate) type BoxedTransactionAction = Arc<dyn TransactionAction>;
+/// A boxed entry pairing a transaction action with its retry-persistent state.
+pub(crate) type TransactionActionEntry = Box<dyn ErasedActionEntry>;
/// A trait representing an atomic action that can be part of a transaction.
///
/// Implementors of this trait define how a specific action is committed to a
table.
/// Each action is responsible for generating the updates and requirements
needed
/// to modify the table metadata.
+///
+/// An action's intent is immutable once applied to a transaction.
Retry-persistent
+/// state lives in the associated [`TransactionAction::State`], which survives
replay
+/// attempts of one logical execution but is never shared between executions
+/// (cloning a transaction creates fresh state via
[`TransactionAction::new_state`]).
#[async_trait]
-pub(crate) trait TransactionAction: AsAny + Sync + Send {
+pub(crate) trait TransactionAction: Clone + Send + Sync + 'static {
+ /// Retry-persistent state exclusively owned by one logical execution of
this
+ /// action. Stateless actions use `State = ()`.
+ type State: Send + Sync + 'static;
+
+ /// Creates fresh state for one logical execution.
+ ///
+ /// This is infallible and table-independent; table-dependent
initialization
+ /// happens during [`TransactionAction::commit`].
+ fn new_state(&self) -> Self::State;
+
/// Commits this action against the provided table and returns the
resulting updates.
/// NOTE: This function is intended for internal use only and should not
be called directly by users.
///
+ /// One replay attempt: the action (intent) is borrowed immutably, its
execution
+ /// state mutably, and the table reflects the current transaction-local
base.
+ ///
/// # Arguments
///
+ /// * `state` - The retry-persistent state for this logical execution.
/// * `table` - The current state of the table this action should apply to.
///
/// # Returns
///
/// An `ActionCommit` containing table updates and table requirements,
/// or an error if the commit fails.
- async fn commit(self: Arc<Self>, table: &Table) -> Result<ActionCommit>;
+ async fn commit(&self, state: &mut Self::State, table: &Table) ->
Result<ActionCommit>;
+
+ /// Best-effort cleanup after the transaction reaches a terminal result.
+ ///
+ /// Consumes the action and its state: the type system guarantees that no
+ /// further attempt can run for this entry after cleanup. Cleanup must not
+ /// change the already-determined transaction result.
+ ///
+ /// See [`CommitStatus`] for what each terminal status allows cleanup to
do.
+ // TODO: invoke this from `Transaction::commit` once terminal status
+ // classification is wired up (stateful transaction RFC, milestone 3).
+ #[allow(dead_code)]
+ async fn cleanup(self: Box<Self>, state: Self::State, table: &Table,
status: CommitStatus);
+}
+
+/// Classification of a transaction's terminal result, consumed by
+/// [`TransactionAction::cleanup`] to determine what is safe to delete.
+///
+/// The transaction/catalog layer determines this classification; actions
+/// consume it rather than independently interpreting catalog errors.
+// TODO: produce this classification in `Transaction::commit` once terminal
+// cleanup is wired up (stateful transaction RFC, milestone 3).
+#[allow(dead_code)]
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub(crate) enum CommitStatus {
+ /// The catalog confirmed that the commit was applied. Cleanup may remove
+ /// owned artifacts not retained by the committed result.
+ Committed,
+ /// The transaction definitively did not commit and will not retry.
+ /// Owned artifacts are deletable.
+ Failed,
+ /// The commit request was submitted, but its outcome could not be
+ /// resolved: the catalog may or may not have applied it. Cleanup must
+ /// delete nothing.
Review Comment:
While we're here — the `Unknown` doc says cleanup must delete nothing, but
the other half of the guarantee is that retry must be suppressed: if a
submitted-but-unconfirmed commit is retried, the original may have landed and
the second attempt duplicates it. Worth stating here so whoever wires milestone
3 classifies submitted-but-unverifiable errors as `Unknown` and keeps them out
of the retry predicate (`when(|e| e.retryable())` has no such exclusion today).
##########
crates/iceberg/src/transaction/action.rs:
##########
@@ -25,28 +24,147 @@ use crate::table::Table;
use crate::transaction::Transaction;
use crate::{Result, TableRequirement, TableUpdate};
-/// A boxed, thread-safe reference to a `TransactionAction`.
-pub(crate) type BoxedTransactionAction = Arc<dyn TransactionAction>;
+/// A boxed entry pairing a transaction action with its retry-persistent state.
+pub(crate) type TransactionActionEntry = Box<dyn ErasedActionEntry>;
/// A trait representing an atomic action that can be part of a transaction.
///
/// Implementors of this trait define how a specific action is committed to a
table.
/// Each action is responsible for generating the updates and requirements
needed
/// to modify the table metadata.
+///
+/// An action's intent is immutable once applied to a transaction.
Retry-persistent
+/// state lives in the associated [`TransactionAction::State`], which survives
replay
+/// attempts of one logical execution but is never shared between executions
+/// (cloning a transaction creates fresh state via
[`TransactionAction::new_state`]).
#[async_trait]
-pub(crate) trait TransactionAction: AsAny + Sync + Send {
+pub(crate) trait TransactionAction: Clone + Send + Sync + 'static {
+ /// Retry-persistent state exclusively owned by one logical execution of
this
+ /// action. Stateless actions use `State = ()`.
+ type State: Send + Sync + 'static;
+
+ /// Creates fresh state for one logical execution.
+ ///
+ /// This is infallible and table-independent; table-dependent
initialization
+ /// happens during [`TransactionAction::commit`].
+ fn new_state(&self) -> Self::State;
+
/// Commits this action against the provided table and returns the
resulting updates.
/// NOTE: This function is intended for internal use only and should not
be called directly by users.
///
+ /// One replay attempt: the action (intent) is borrowed immutably, its
execution
+ /// state mutably, and the table reflects the current transaction-local
base.
+ ///
/// # Arguments
///
+ /// * `state` - The retry-persistent state for this logical execution.
/// * `table` - The current state of the table this action should apply to.
///
/// # Returns
///
/// An `ActionCommit` containing table updates and table requirements,
/// or an error if the commit fails.
- async fn commit(self: Arc<Self>, table: &Table) -> Result<ActionCommit>;
+ async fn commit(&self, state: &mut Self::State, table: &Table) ->
Result<ActionCommit>;
+
+ /// Best-effort cleanup after the transaction reaches a terminal result.
+ ///
+ /// Consumes the action and its state: the type system guarantees that no
+ /// further attempt can run for this entry after cleanup. Cleanup must not
+ /// change the already-determined transaction result.
+ ///
+ /// See [`CommitStatus`] for what each terminal status allows cleanup to
do.
+ // TODO: invoke this from `Transaction::commit` once terminal status
+ // classification is wired up (stateful transaction RFC, milestone 3).
+ #[allow(dead_code)]
+ async fn cleanup(self: Box<Self>, state: Self::State, table: &Table,
status: CommitStatus);
+}
+
+/// Classification of a transaction's terminal result, consumed by
+/// [`TransactionAction::cleanup`] to determine what is safe to delete.
+///
+/// The transaction/catalog layer determines this classification; actions
+/// consume it rather than independently interpreting catalog errors.
+// TODO: produce this classification in `Transaction::commit` once terminal
+// cleanup is wired up (stateful transaction RFC, milestone 3).
+#[allow(dead_code)]
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub(crate) enum CommitStatus {
+ /// The catalog confirmed that the commit was applied. Cleanup may remove
+ /// owned artifacts not retained by the committed result.
+ Committed,
+ /// The transaction definitively did not commit and will not retry.
+ /// Owned artifacts are deletable.
+ Failed,
+ /// The commit request was submitted, but its outcome could not be
+ /// resolved: the catalog may or may not have applied it. Cleanup must
+ /// delete nothing.
+ ///
+ /// Example: every action executed and validated successfully, but the
+ /// connection failed while awaiting the catalog's response to
+ /// `update_table`.
+ Unknown,
+}
+
+/// An entry in a transaction, pairing an action's immutable intent with the
+/// retry-persistent state of one logical execution.
+///
+/// The pairing is preserved by construction: the entry is created with fresh
+/// state and owns both exclusively, so terminal cleanup can consume them
together.
+pub(crate) struct ActionEntry<A: TransactionAction> {
+ action: Box<A>,
Review Comment:
`ActionEntry<A>` stores `action: Box<A>` but is itself always held behind
`Box<dyn ErasedActionEntry>`, so every entry is two allocations and the hot
`do_commit` loop derefs through the inner box. The inner box only earns its
keep at `cleanup`, which needs `self: Box<Self>` on the action. I'd store
`action: A` and box once at that cold cleanup call
(`Box::new(entry.action).cleanup(...)`); if the pre-boxing is deliberate, a
one-line comment saying why would stop the next reader from "simplifying" it
and breaking cleanup forwarding.
##########
crates/iceberg/src/transaction/action.rs:
##########
@@ -159,14 +298,30 @@ mod tests {
let tx = Transaction::new(&table);
let updated_tx = action.apply(tx).unwrap();
- // There should be one action in the transaction now
+ // There should be one action entry in the transaction now
assert_eq!(updated_tx.actions.len(), 1);
(*updated_tx.actions[0])
- .downcast_ref::<TestAction>()
+ .downcast_ref::<ActionEntry<TestAction>>()
.expect("TestAction was not applied to Transaction!");
}
+ #[test]
+ fn test_transaction_clone_creates_fresh_entries() {
Review Comment:
This test can't verify the invariant it names: with `TestAction::State =
()`, a fresh-state clone and a state-sharing clone are indistinguishable, so
the milestone's central claim — retained state across retries, fresh state
across clones — is exercised in neither direction.
I'd add a test action with non-trivial state (say `State = u32` bumped on
each `commit`) and assert both halves: `fresh_clone` resets the counter to
zero, and driving `do_commit` twice on the same entry accumulates. That's the
one thing this milestone most needs to pin down.
##########
crates/iceberg/src/transaction/action.rs:
##########
@@ -25,28 +24,147 @@ use crate::table::Table;
use crate::transaction::Transaction;
use crate::{Result, TableRequirement, TableUpdate};
-/// A boxed, thread-safe reference to a `TransactionAction`.
-pub(crate) type BoxedTransactionAction = Arc<dyn TransactionAction>;
+/// A boxed entry pairing a transaction action with its retry-persistent state.
+pub(crate) type TransactionActionEntry = Box<dyn ErasedActionEntry>;
/// A trait representing an atomic action that can be part of a transaction.
///
/// Implementors of this trait define how a specific action is committed to a
table.
/// Each action is responsible for generating the updates and requirements
needed
/// to modify the table metadata.
+///
+/// An action's intent is immutable once applied to a transaction.
Retry-persistent
+/// state lives in the associated [`TransactionAction::State`], which survives
replay
+/// attempts of one logical execution but is never shared between executions
+/// (cloning a transaction creates fresh state via
[`TransactionAction::new_state`]).
#[async_trait]
-pub(crate) trait TransactionAction: AsAny + Sync + Send {
+pub(crate) trait TransactionAction: Clone + Send + Sync + 'static {
+ /// Retry-persistent state exclusively owned by one logical execution of
this
+ /// action. Stateless actions use `State = ()`.
+ type State: Send + Sync + 'static;
+
+ /// Creates fresh state for one logical execution.
+ ///
+ /// This is infallible and table-independent; table-dependent
initialization
+ /// happens during [`TransactionAction::commit`].
+ fn new_state(&self) -> Self::State;
+
/// Commits this action against the provided table and returns the
resulting updates.
/// NOTE: This function is intended for internal use only and should not
be called directly by users.
///
+ /// One replay attempt: the action (intent) is borrowed immutably, its
execution
+ /// state mutably, and the table reflects the current transaction-local
base.
+ ///
/// # Arguments
///
+ /// * `state` - The retry-persistent state for this logical execution.
/// * `table` - The current state of the table this action should apply to.
///
/// # Returns
///
/// An `ActionCommit` containing table updates and table requirements,
/// or an error if the commit fails.
- async fn commit(self: Arc<Self>, table: &Table) -> Result<ActionCommit>;
+ async fn commit(&self, state: &mut Self::State, table: &Table) ->
Result<ActionCommit>;
+
+ /// Best-effort cleanup after the transaction reaches a terminal result.
+ ///
+ /// Consumes the action and its state: the type system guarantees that no
+ /// further attempt can run for this entry after cleanup. Cleanup must not
+ /// change the already-determined transaction result.
+ ///
+ /// See [`CommitStatus`] for what each terminal status allows cleanup to
do.
+ // TODO: invoke this from `Transaction::commit` once terminal status
+ // classification is wired up (stateful transaction RFC, milestone 3).
+ #[allow(dead_code)]
+ async fn cleanup(self: Box<Self>, state: Self::State, table: &Table,
status: CommitStatus);
+}
+
+/// Classification of a transaction's terminal result, consumed by
+/// [`TransactionAction::cleanup`] to determine what is safe to delete.
+///
+/// The transaction/catalog layer determines this classification; actions
+/// consume it rather than independently interpreting catalog errors.
+// TODO: produce this classification in `Transaction::commit` once terminal
+// cleanup is wired up (stateful transaction RFC, milestone 3).
+#[allow(dead_code)]
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub(crate) enum CommitStatus {
+ /// The catalog confirmed that the commit was applied. Cleanup may remove
+ /// owned artifacts not retained by the committed result.
+ Committed,
+ /// The transaction definitively did not commit and will not retry.
+ /// Owned artifacts are deletable.
+ Failed,
+ /// The commit request was submitted, but its outcome could not be
+ /// resolved: the catalog may or may not have applied it. Cleanup must
+ /// delete nothing.
+ ///
+ /// Example: every action executed and validated successfully, but the
+ /// connection failed while awaiting the catalog's response to
+ /// `update_table`.
+ Unknown,
+}
+
+/// An entry in a transaction, pairing an action's immutable intent with the
+/// retry-persistent state of one logical execution.
+///
+/// The pairing is preserved by construction: the entry is created with fresh
+/// state and owns both exclusively, so terminal cleanup can consume them
together.
+pub(crate) struct ActionEntry<A: TransactionAction> {
+ action: Box<A>,
+ state: A::State,
+}
+
+impl<A: TransactionAction> ActionEntry<A> {
+ fn new(action: A) -> Self {
+ let state = action.new_state();
+ Self {
+ action: Box::new(action),
+ state,
+ }
+ }
+
+ /// The action (intent) held by this entry.
+ #[cfg(test)]
+ pub(crate) fn action(&self) -> &A {
+ &self.action
+ }
+}
+
+/// Object-safe adapter over [`ActionEntry`], allowing a transaction to store
+/// heterogeneous entries while keeping each action paired with its own state.
+#[async_trait]
+pub(crate) trait ErasedActionEntry: AsAny + Send + Sync {
+ /// One replay attempt against the current transaction-local table.
+ async fn commit(&mut self, table: &Table) -> Result<ActionCommit>;
+
+ /// Best-effort terminal cleanup, consuming the entry.
+ // TODO: invoke from `Transaction::commit` once terminal status
+ // classification is wired up (stateful transaction RFC, milestone 3).
+ #[allow(dead_code)]
+ async fn cleanup(self: Box<Self>, table: &Table, status: CommitStatus);
+
+ /// Creates a new entry for a new logical execution of the same action:
+ /// the action (intent) is cloned, and it is paired with fresh state from
+ /// [`TransactionAction::new_state`]. Accrued retry state is never carried
+ /// over to the new entry.
+ fn fresh_clone(&self) -> TransactionActionEntry;
+}
+
+#[async_trait]
+impl<A: TransactionAction> ErasedActionEntry for ActionEntry<A> {
+ async fn commit(&mut self, table: &Table) -> Result<ActionCommit> {
+ self.action.commit(&mut self.state, table).await
+ }
+
+ async fn cleanup(self: Box<Self>, table: &Table, status: CommitStatus) {
+ let entry = *self;
Review Comment:
This unbox-and-forward path is the most novel code in the PR and nothing
exercises it — it's `#[allow(dead_code)]` and no test constructs a
`CommitStatus` or calls `cleanup`. I'd add a two-line unit test that boxes an
`ActionEntry` and calls `.cleanup(&table, CommitStatus::Committed).await`, so
the forwarding chain is validated now rather than discovered when milestone 3
wires it up.
##########
crates/iceberg/src/transaction/action.rs:
##########
@@ -25,28 +24,147 @@ use crate::table::Table;
use crate::transaction::Transaction;
use crate::{Result, TableRequirement, TableUpdate};
-/// A boxed, thread-safe reference to a `TransactionAction`.
-pub(crate) type BoxedTransactionAction = Arc<dyn TransactionAction>;
+/// A boxed entry pairing a transaction action with its retry-persistent state.
+pub(crate) type TransactionActionEntry = Box<dyn ErasedActionEntry>;
/// A trait representing an atomic action that can be part of a transaction.
///
/// Implementors of this trait define how a specific action is committed to a
table.
/// Each action is responsible for generating the updates and requirements
needed
/// to modify the table metadata.
+///
+/// An action's intent is immutable once applied to a transaction.
Retry-persistent
+/// state lives in the associated [`TransactionAction::State`], which survives
replay
+/// attempts of one logical execution but is never shared between executions
+/// (cloning a transaction creates fresh state via
[`TransactionAction::new_state`]).
#[async_trait]
-pub(crate) trait TransactionAction: AsAny + Sync + Send {
+pub(crate) trait TransactionAction: Clone + Send + Sync + 'static {
+ /// Retry-persistent state exclusively owned by one logical execution of
this
+ /// action. Stateless actions use `State = ()`.
+ type State: Send + Sync + 'static;
+
+ /// Creates fresh state for one logical execution.
+ ///
+ /// This is infallible and table-independent; table-dependent
initialization
+ /// happens during [`TransactionAction::commit`].
+ fn new_state(&self) -> Self::State;
+
/// Commits this action against the provided table and returns the
resulting updates.
/// NOTE: This function is intended for internal use only and should not
be called directly by users.
///
+ /// One replay attempt: the action (intent) is borrowed immutably, its
execution
+ /// state mutably, and the table reflects the current transaction-local
base.
+ ///
/// # Arguments
///
+ /// * `state` - The retry-persistent state for this logical execution.
/// * `table` - The current state of the table this action should apply to.
///
/// # Returns
///
/// An `ActionCommit` containing table updates and table requirements,
/// or an error if the commit fails.
- async fn commit(self: Arc<Self>, table: &Table) -> Result<ActionCommit>;
+ async fn commit(&self, state: &mut Self::State, table: &Table) ->
Result<ActionCommit>;
+
+ /// Best-effort cleanup after the transaction reaches a terminal result.
+ ///
+ /// Consumes the action and its state: the type system guarantees that no
+ /// further attempt can run for this entry after cleanup. Cleanup must not
+ /// change the already-determined transaction result.
+ ///
+ /// See [`CommitStatus`] for what each terminal status allows cleanup to
do.
+ // TODO: invoke this from `Transaction::commit` once terminal status
+ // classification is wired up (stateful transaction RFC, milestone 3).
+ #[allow(dead_code)]
+ async fn cleanup(self: Box<Self>, state: Self::State, table: &Table,
status: CommitStatus);
+}
+
+/// Classification of a transaction's terminal result, consumed by
+/// [`TransactionAction::cleanup`] to determine what is safe to delete.
+///
+/// The transaction/catalog layer determines this classification; actions
+/// consume it rather than independently interpreting catalog errors.
+// TODO: produce this classification in `Transaction::commit` once terminal
+// cleanup is wired up (stateful transaction RFC, milestone 3).
+#[allow(dead_code)]
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub(crate) enum CommitStatus {
+ /// The catalog confirmed that the commit was applied. Cleanup may remove
+ /// owned artifacts not retained by the committed result.
+ Committed,
+ /// The transaction definitively did not commit and will not retry.
+ /// Owned artifacts are deletable.
+ Failed,
+ /// The commit request was submitted, but its outcome could not be
+ /// resolved: the catalog may or may not have applied it. Cleanup must
+ /// delete nothing.
+ ///
+ /// Example: every action executed and validated successfully, but the
+ /// connection failed while awaiting the catalog's response to
+ /// `update_table`.
+ Unknown,
+}
+
+/// An entry in a transaction, pairing an action's immutable intent with the
+/// retry-persistent state of one logical execution.
+///
+/// The pairing is preserved by construction: the entry is created with fresh
+/// state and owns both exclusively, so terminal cleanup can consume them
together.
+pub(crate) struct ActionEntry<A: TransactionAction> {
+ action: Box<A>,
+ state: A::State,
+}
+
+impl<A: TransactionAction> ActionEntry<A> {
+ fn new(action: A) -> Self {
+ let state = action.new_state();
+ Self {
+ action: Box::new(action),
+ state,
+ }
+ }
+
+ /// The action (intent) held by this entry.
+ #[cfg(test)]
+ pub(crate) fn action(&self) -> &A {
+ &self.action
+ }
+}
+
+/// Object-safe adapter over [`ActionEntry`], allowing a transaction to store
+/// heterogeneous entries while keeping each action paired with its own state.
+#[async_trait]
+pub(crate) trait ErasedActionEntry: AsAny + Send + Sync {
+ /// One replay attempt against the current transaction-local table.
+ async fn commit(&mut self, table: &Table) -> Result<ActionCommit>;
+
+ /// Best-effort terminal cleanup, consuming the entry.
+ // TODO: invoke from `Transaction::commit` once terminal status
+ // classification is wired up (stateful transaction RFC, milestone 3).
+ #[allow(dead_code)]
+ async fn cleanup(self: Box<Self>, table: &Table, status: CommitStatus);
+
+ /// Creates a new entry for a new logical execution of the same action:
+ /// the action (intent) is cloned, and it is paired with fresh state from
+ /// [`TransactionAction::new_state`]. Accrued retry state is never carried
+ /// over to the new entry.
+ fn fresh_clone(&self) -> TransactionActionEntry;
Review Comment:
Minor: `fresh_clone` reads like the clone is fresh, but it's the *state*
that's fresh — the action intent is cloned as-is. Something like
`clone_with_fresh_state` or `new_execution_entry` would say what it does. Not
blocking.
##########
crates/iceberg/src/transaction/update_location.rs:
##########
@@ -56,7 +55,11 @@ impl UpdateLocationAction {
#[async_trait]
impl TransactionAction for UpdateLocationAction {
- async fn commit(self: Arc<Self>, _table: &Table) -> Result<ActionCommit> {
+ type State = ();
+
+ fn new_state(&self) -> Self::State {}
+
+ async fn commit(&self, _state: &mut (), _table: &Table) ->
Result<ActionCommit> {
let updates: Vec<TableUpdate>;
if let Some(location) = self.location.clone() {
Review Comment:
Small thing: `self.location.clone()` clones the whole `Option<String>` just
to destructure it. `as_deref` avoids the clone:
```rust
if let Some(location) = self.location.as_deref() {
updates = vec![TableUpdate::SetLocation { location: location.to_owned()
}];
}
```
##########
crates/iceberg/src/transaction/append.rs:
##########
@@ -87,7 +88,13 @@ impl FastAppendAction {
#[async_trait]
impl TransactionAction for FastAppendAction {
- async fn commit(self: Arc<Self>, table: &Table) -> Result<ActionCommit> {
+ // TODO: replace with a persistent `SimpleSnapshotProducer` once snapshot
+ // production is migrated (stateful transaction RFC, milestone 2).
+ type State = ();
+
+ fn new_state(&self) -> Self::State {}
+
+ async fn commit(&self, _state: &mut (), table: &Table) ->
Result<ActionCommit> {
let snapshot_producer = SnapshotProducer::new(
table,
self.commit_uuid.unwrap_or_else(Uuid::now_v7),
Review Comment:
Not for this PR — `commit_uuid` defaults to a fresh `Uuid::now_v7()` on
every `commit` call, so once retries drive `commit` more than once, each
attempt writes manifests under a new UUID and the no-op cleanup never reclaims
the orphans. It's pre-existing behavior, but this is the right spot to make the
milestone-2 TODO explicit: moving the `SnapshotProducer` (and its UUID) into
`new_state` is what makes retries reuse the same UUID and already-written
manifests.
--
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]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]