Explorar o código

Make posting reads deterministic and migrations atomic

Reviewing the append-only value-table / hot-index split surfaced three
follow-ups worth fixing.

Migrations ran statement-by-statement outside any transaction, recording
a migration as applied only after all of its statements succeeded. The
new migration that drops and rebuilds the postings table could therefore
crash half-applied, leaving the schema in a state the migration could not
be re-run against. Wrap each migration's statements and its bookkeeping
insert in one transaction so a crash rolls back cleanly and the migration
is retried whole. Both SQLite and PostgreSQL support transactional DDL.

Posting reads returned rows in an unspecified order, so LIMIT/OFFSET
pagination could skip or repeat rows across pages, worst for the live set
whose source is a UNION of two tables. Order get_postings_by_account by
the posting id in both backends so the primitive that balance, selection,
close, and pagination all build on returns a stable sequence.

Bring the reference docs back in line with the derived-state model:
postings are immutable and carry no lifecycle column; state is Active,
Reserved, Spent, or Missing derived from index membership.
Cesar Rodas hai 1 semana
pai
achega
41ce95cb4d

+ 29 - 3
crates/kuatia-storage-sql/src/lib.rs

@@ -125,11 +125,23 @@ impl SqlStore {
                 continue;
                 continue;
             }
             }
 
 
+            // Apply every statement and record the migration in one transaction,
+            // so a crash mid-migration rolls back cleanly and the migration is
+            // retried as a whole. Migration 004 drops and rebuilds `postings`;
+            // without the transaction a partial apply would leave the schema in a
+            // state the migration cannot be re-run against. Both SQLite and
+            // PostgreSQL support transactional DDL.
+            let mut tx = self
+                .pool
+                .begin()
+                .await
+                .map_err(|e| StoreError::Internal(e.to_string()))?;
+
             for statement in sql.split(';') {
             for statement in sql.split(';') {
                 let trimmed = statement.trim();
                 let trimmed = statement.trim();
                 if !trimmed.is_empty() {
                 if !trimmed.is_empty() {
                     sqlx::query(trimmed)
                     sqlx::query(trimmed)
-                        .execute(&self.pool)
+                        .execute(&mut *tx)
                         .await
                         .await
                         .map_err(|e| StoreError::Internal(e.to_string()))?;
                         .map_err(|e| StoreError::Internal(e.to_string()))?;
                 }
                 }
@@ -137,7 +149,11 @@ impl SqlStore {
 
 
             sqlx::query("INSERT INTO _migrations (name) VALUES ($1)")
             sqlx::query("INSERT INTO _migrations (name) VALUES ($1)")
                 .bind(*name)
                 .bind(*name)
-                .execute(&self.pool)
+                .execute(&mut *tx)
+                .await
+                .map_err(|e| StoreError::Internal(e.to_string()))?;
+
+            tx.commit()
                 .await
                 .await
                 .map_err(|e| StoreError::Internal(e.to_string()))?;
                 .map_err(|e| StoreError::Internal(e.to_string()))?;
         }
         }
@@ -586,6 +602,10 @@ impl PostingStore for SqlStore {
         if asset.is_some() {
         if asset.is_some() {
             sql.push_str(&format!(" AND asset = ${placeholder}"));
             sql.push_str(&format!(" AND asset = ${placeholder}"));
         }
         }
+        // Deterministic order by the posting primary key, matching
+        // `query_postings`, so callers (and pagination built on top) see a
+        // stable sequence.
+        sql.push_str(" ORDER BY transfer_id, idx");
 
 
         let mut q = sqlx::query(&sql).bind(id);
         let mut q = sqlx::query(&sql).bind(id);
         if let Some(s) = sub {
         if let Some(s) = sub {
@@ -712,7 +732,13 @@ impl PostingStore for SqlStore {
             let c = format!("SELECT COUNT(*) as cnt FROM {source} {w}");
             let c = format!("SELECT COUNT(*) as cnt FROM {source} {w}");
             let limit = query.limit.unwrap_or(u32::MAX);
             let limit = query.limit.unwrap_or(u32::MAX);
             let offset = query.offset.unwrap_or(0);
             let offset = query.offset.unwrap_or(0);
-            w.push_str(&format!(" LIMIT {limit} OFFSET {offset}"));
+            // Order by the posting primary key so pagination is deterministic:
+            // without it LIMIT/OFFSET could skip or repeat rows across pages,
+            // especially for `Live`, whose source is a `UNION ALL` with no
+            // inherent order.
+            w.push_str(&format!(
+                " ORDER BY transfer_id, idx LIMIT {limit} OFFSET {offset}"
+            ));
             (format!("SELECT * FROM {source} {w}"), c)
             (format!("SELECT * FROM {source} {w}"), c)
         };
         };
 
 

+ 4 - 1
crates/kuatia-storage/src/mem_store.rs

@@ -176,7 +176,7 @@ impl PostingStore for InMemoryStore {
         // Read from the relevant set directly — no merge back to the immutable
         // Read from the relevant set directly — no merge back to the immutable
         // record. `Live` is the two disjoint live sets chained together.
         // record. `Live` is the two disjoint live sets chained together.
         let reserved = || postings.reserved.values().map(|(p, _)| p);
         let reserved = || postings.reserved.values().map(|(p, _)| p);
-        let result: Vec<Posting> = match filter {
+        let mut result: Vec<Posting> = match filter {
             PostingFilter::Active => postings.active.values().filter(matches).cloned().collect(),
             PostingFilter::Active => postings.active.values().filter(matches).cloned().collect(),
             PostingFilter::Reserved => reserved().filter(matches).cloned().collect(),
             PostingFilter::Reserved => reserved().filter(matches).cloned().collect(),
             PostingFilter::Live => postings
             PostingFilter::Live => postings
@@ -193,6 +193,9 @@ impl PostingStore for InMemoryStore {
                 .cloned()
                 .cloned()
                 .collect(),
                 .collect(),
         };
         };
+        // Deterministic order by the posting id, matching the SQL backend and
+        // `query_postings`; the maps above iterate in an unspecified order.
+        result.sort_by_key(|p| p.id);
         Ok(result)
         Ok(result)
     }
     }
 
 

+ 5 - 1
crates/kuatia-storage/src/store.rs

@@ -111,7 +111,9 @@ pub trait PostingStore: Send + Sync {
     async fn get_postings(&self, ids: &[PostingId]) -> Result<Vec<Posting>, StoreError>;
     async fn get_postings(&self, ids: &[PostingId]) -> Result<Vec<Posting>, StoreError>;
     /// Return postings owned by a base account, filtered by subaccount, asset,
     /// Return postings owned by a base account, filtered by subaccount, asset,
     /// and derived lifecycle state (see [`PostingFilter`]). `sub == None` spans
     /// and derived lifecycle state (see [`PostingFilter`]). `sub == None` spans
-    /// every subaccount of `id`; `sub == Some(s)` restricts to that one.
+    /// every subaccount of `id`; `sub == Some(s)` restricts to that one. Results
+    /// are ordered by posting id, so callers and pagination built on top see a
+    /// stable sequence.
     async fn get_postings_by_account(
     async fn get_postings_by_account(
         &self,
         &self,
         id: i64,
         id: i64,
@@ -167,6 +169,8 @@ pub trait PostingStore: Send + Sync {
 
 
     /// Query postings with filtering and pagination.
     /// Query postings with filtering and pagination.
     async fn query_postings(&self, query: &PostingQuery) -> Result<Page<Posting>, StoreError> {
     async fn query_postings(&self, query: &PostingQuery) -> Result<Page<Posting>, StoreError> {
+        // `get_postings_by_account` returns a deterministic id-ordered sequence,
+        // so LIMIT/OFFSET here paginate stably without an extra sort.
         let all = self
         let all = self
             .get_postings_by_account(query.account, query.sub, query.asset.as_ref(), query.filter)
             .get_postings_by_account(query.account, query.sub, query.asset.as_ref(), query.filter)
             .await?;
             .await?;

+ 3 - 4
doc/accounts.md

@@ -4,10 +4,9 @@
 
 
 An account is a versioned entity that owns postings. Balance is never
 An account is a versioned entity that owns postings. Balance is never
 stored: it is always computed from postings for a given (account, asset)
 stored: it is always computed from postings for a given (account, asset)
-pair. The ledger balance sums non-`Inactive` postings (`Active +
-PendingInactive`); the available balance sums only `Active` postings
-(excluding those reserved for an in-flight transfer). `balance()` returns
-the ledger balance.
+pair. `balance()` sums the live postings, meaning those in the active or
+reserved index (`Active ∪ Reserved`); spent postings, which remain only in
+the immutable table, are excluded.
 
 
 ## Structure
 ## Structure
 
 

+ 30 - 25
doc/crates.md

@@ -25,10 +25,11 @@ dependencies (`sha2`, `serde`, `bitflags`).
 | `AccountSnapshotId { account, snapshot_id }` | Account state hash for version pinning |
 | `AccountSnapshotId { account, snapshot_id }` | Account state hash for version pinning |
 | `Cent` | Smallest monetary unit (private field, backing integer hidden). Backing is `i64` by default, `i128` under the `i128` feature. Checked arithmetic via `checked_add`, `checked_sub`, `checked_neg`, `checked_sum` returning `Result<Cent, OverflowError>` |
 | `Cent` | Smallest monetary unit (private field, backing integer hidden). Backing is `i64` by default, `i128` under the `i128` feature. Checked arithmetic via `checked_add`, `checked_sub`, `checked_neg`, `checked_sum` returning `Result<Cent, OverflowError>` |
 | `OverflowError` | Returned when a `Cent` operation would overflow or underflow |
 | `OverflowError` | Returned when a `Cent` operation would overflow or underflow |
-| `PostingStatus` | Posting lifecycle: `Active`, `PendingInactive`, `Inactive` |
+| `PostingState` | Derived posting lifecycle (from index-table membership, not stored): `Active`, `Reserved(ReservationId)`, `Spent`, `Missing` |
+| `PostingFilter` | Read filter over derived posting state: `Active`, `Reserved`, `Live` (Active ∪ Reserved), `All` |
 | `Amount` | Parser/formatter for decimal strings. Not stored; use at API boundaries only |
 | `Amount` | Parser/formatter for decimal strings. Not stored; use at API boundaries only |
-| `Posting` | Signed amount of one asset owned by one account. Has `status: PostingStatus` and `reservation: Option<ReservationId>` (owner token while `PendingInactive`) |
-| `ReservationId` | Owner token stamped on a reserved posting so only the reserving saga may finalize/release it |
+| `Posting` | Immutable signed amount of one asset owned by one account. Carries no lifecycle field; its state is derived from index-table membership (see `PostingState`) |
+| `ReservationId` | Owner token recorded in the reserved index so only the reserving saga may finalize/release a posting |
 | `NewPosting` | Posting to be created (no id yet, assigned during validation) |
 | `NewPosting` | Posting to be created (no id yet, assigned during validation) |
 | `Transfer` | Atomic unit: consumes postings + creates postings + metadata |
 | `Transfer` | Atomic unit: consumes postings + creates postings + metadata |
 | `EnvelopeBuilder` | Fluent builder for `Transfer` construction |
 | `EnvelopeBuilder` | Fluent builder for `Transfer` construction |
@@ -48,8 +49,7 @@ in order:
 graph TD
 graph TD
     A[1. Non-empty] --> B[2. No duplicate consumes]
     A[1. Non-empty] --> B[2. No duplicate consumes]
     B --> C[3. Posting existence]
     B --> C[3. Posting existence]
-    C --> D[4. Posting active or reserved]
-    D --> E[5. Account existence & lifecycle]
+    C --> E[5. Account existence & lifecycle]
     E --> F[6. Snapshot pinning]
     E --> F[6. Snapshot pinning]
     F --> BP[7. Book policy]
     F --> BP[7. Book policy]
     BP --> G[8. Per-asset conservation]
     BP --> G[8. Per-asset conservation]
@@ -61,9 +61,9 @@ graph TD
 
 
 1. **Non-empty**: transfer must consume or create at least one posting
 1. **Non-empty**: transfer must consume or create at least one posting
 2. **No duplicate consumes**: each posting consumed at most once
 2. **No duplicate consumes**: each posting consumed at most once
-3. **Posting existence**: every consumed posting exists in state
-4. **Posting active or reserved**: consumed postings must be `Active` or
-   `PendingInactive` (prevents double-spend)
+3. **Posting existence**: every consumed posting exists in the immutable table
+   (a `Posting` carries no lifecycle state; double-spend safety is enforced by
+   the reserve claim and the finalize "all spent" guard, not here)
 5. **Account existence & lifecycle**: all referenced accounts exist, not
 5. **Account existence & lifecycle**: all referenced accounts exist, not
    frozen, not closed
    frozen, not closed
 6. **Snapshot pinning**: account snapshots (if provided) must match current
 6. **Snapshot pinning**: account snapshots (if provided) must match current
@@ -108,7 +108,7 @@ writes, then calls the dumb primitives, interpreting/verifying each count:
 graph LR
 graph LR
     A[resolve] -->|Envelope| W[save PendingSaga: Reserving]
     A[resolve] -->|Envelope| W[save PendingSaga: Reserving]
     W --> B[reserve_postings]
     W --> B[reserve_postings]
-    B -->|Active→PendingInactive| F[finalize]
+    B -->|active index → reserved index| F[finalize]
     F --> V[validate_and_plan re-check]
     F --> V[validate_and_plan re-check]
     V --> M[mark Finalizing]
     V --> M[mark Finalizing]
     M --> D[deactivate → insert → store_transfer → append_event]
     M --> D[deactivate → insert → store_transfer → append_event]
@@ -157,13 +157,13 @@ Transfers are built via `TransferBuilder` and committed with
 
 
 | Method | Description |
 | Method | Description |
 |--------|-------------|
 |--------|-------------|
-| `balance(account, asset)` | Sum of non-Inactive postings (computed by Ledger) |
+| `balance(account, asset)` | Sum of live (Active or Reserved) postings (computed by Ledger) |
 | `list_accounts()` | All current account snapshots |
 | `list_accounts()` | All current account snapshots |
 | `get_account(id)` | Latest account snapshot |
 | `get_account(id)` | Latest account snapshot |
 | `query_transfers(query)` | Paginated, filtered transfer history (by date range, book) |
 | `query_transfers(query)` | Paginated, filtered transfer history (by date range, book) |
 | `history(account)` | All transfers involving an account |
 | `history(account)` | All transfers involving an account |
-| `postings(account)` | All postings (any status) |
-| `query_postings(query)` | Paginated, filtered postings (by asset, status) |
+| `postings(account)` | All postings (any state) |
+| `query_postings(query)` | Paginated, filtered postings (by asset, `PostingFilter`) |
 | `account_history(id)` | All version snapshots |
 | `account_history(id)` | All version snapshots |
 | `get_events_since(seq, limit)` | Query ledger event log after a sequence number |
 | `get_events_since(seq, limit)` | Query ledger event log after a sequence number |
 
 
@@ -185,8 +185,8 @@ graph TB
 
 
 - **`AccountStore`**: `get_account`, `get_accounts`, `create_account`,
 - **`AccountStore`**: `get_account`, `get_accounts`, `create_account`,
   `append_account_version`, `get_account_history`, `list_accounts`
   `append_account_version`, `get_account_history`, `list_accounts`
-- **`PostingStore`**: `get_postings`,
-  `get_postings_by_account(account, asset?, status?)`, `query_postings(query)`,
+- **`PostingStore`**: `get_postings`, `get_posting_states`,
+  `get_postings_by_account(account, sub?, asset?, filter)`, `query_postings(query)`,
   and the dumb write primitives `reserve_postings(ids, reservation) -> u64`,
   and the dumb write primitives `reserve_postings(ids, reservation) -> u64`,
   `release_postings(ids, reservation) -> u64`,
   `release_postings(ids, reservation) -> u64`,
   `deactivate_postings(ids, reservation?) -> u64`,
   `deactivate_postings(ids, reservation?) -> u64`,
@@ -210,24 +210,29 @@ recovery rather than a single transaction.
 conditional update and return how many rows changed (the saga decides what a
 conditional update and return how many rows changed (the saga decides what a
 short count means):
 short count means):
 
 
+State is derived from which index a posting is in: active index → `Active`,
+reserved index → `Reserved(rid)`, neither (only the immutable table) → `Spent`.
+Every transition is an insert/delete on an index table; the posting row never
+changes.
+
 ```mermaid
 ```mermaid
 stateDiagram-v2
 stateDiagram-v2
     [*] --> Active: insert_postings
     [*] --> Active: insert_postings
-    Active --> PendingInactive: reserve_postings
-    PendingInactive --> Active: release_postings
-    PendingInactive --> Inactive: deactivate_postings(reservation)
-    Active --> Inactive: deactivate_postings(None)
+    Active --> Reserved: reserve_postings
+    Reserved --> Active: release_postings
+    Reserved --> Spent: deactivate_postings(reservation)
+    Active --> Spent: deactivate_postings(None)
 ```
 ```
 
 
-Each cell is the count a primitive returns (1 = flipped, 0 = no-op / not
+Each cell is the count a primitive returns (1 = moved, 0 = no-op / not
 applicable). The saga interprets a 0:
 applicable). The saga interprets a 0:
 
 
-| Operation | Active | PendingInactive (this rid) | Inactive |
-|-----------|--------|----------------------------|----------|
-| `reserve_postings(rid)` | → PendingInactive (1) | 0 | 0 |
+| Operation | Active | Reserved (this rid) | Spent |
+|-----------|--------|---------------------|-------|
+| `reserve_postings(rid)` | → Reserved (1) | 0 | 0 |
 | `release_postings(rid)` | 0 | → Active (1) | 0 |
 | `release_postings(rid)` | 0 | → Active (1) | 0 |
-| `deactivate_postings(Some rid)` | 0 | → Inactive (1) | 0 |
-| `deactivate_postings(None)` | → Inactive (1) | 0 | 0 |
+| `deactivate_postings(Some rid)` | 0 | → Spent (1) | 0 |
+| `deactivate_postings(None)` | → Spent (1) | 0 | 0 |
 
 
 There is no all-or-nothing batch rejection: a posting whose condition does not
 There is no all-or-nothing batch rejection: a posting whose condition does not
 hold is skipped (counted as 0, not an error), so a call can apply to some ids
 hold is skipped (counted as 0, not an error), so a call can apply to some ids
@@ -271,7 +276,7 @@ the saga derives meaning from them.
 
 
 | Step | Execute | Compensate | Retry |
 | Step | Execute | Compensate | Retry |
 |------|---------|------------|-------|
 |------|---------|------------|-------|
-| `ReservePostingsStep` | `reserve_postings` `Active → PendingInactive`, interpret count | Release back to `Active` | 3 |
+| `ReservePostingsStep` | `reserve_postings`: active index → reserved index, interpret count | Release back to `Active` | 3 |
 | `FinalizeTransferStep` | `Ledger::finalize_envelope`: re-validate (last-step floor/freeze guard) → mark `Finalizing` → `deactivate` → `insert` → `store_transfer` → `append_event`, verifying every end-state | `reverse(transfer_id)` | 3 |
 | `FinalizeTransferStep` | `Ledger::finalize_envelope`: re-validate (last-step floor/freeze guard) → mark `Finalizing` → `deactivate` → `insert` → `store_transfer` → `append_event`, verifying every end-state | `reverse(transfer_id)` | 3 |
 
 
 Validation lives inside the finalize step so it runs immediately before the
 Validation lives inside the finalize step so it runs immediately before the

+ 3 - 2
doc/transfers.md

@@ -216,8 +216,9 @@ validation steps are:
 
 
 1. Non-empty (must consume or create at least one posting)
 1. Non-empty (must consume or create at least one posting)
 2. No duplicate consumed PostingIds
 2. No duplicate consumed PostingIds
-3. All consumed postings exist
-4. All consumed postings are Active or PendingInactive
+3. All consumed postings exist in the immutable table (a `Posting` carries no
+   lifecycle state; double-spend safety is enforced by the reserve claim and the
+   finalize "all spent" guard, not by validation)
 5. All referenced accounts exist, not frozen, not closed
 5. All referenced accounts exist, not frozen, not closed
 6. Account snapshot pinning (if provided)
 6. Account snapshot pinning (if provided)
 7. Book policy (if a book is loaded): referenced assets/accounts/flags allowed
 7. Book policy (if a book is loaded): referenced assets/accounts/flags allowed