Commit fa860fd
authored
feat: wire outbox worker in financial-accounting service (#422)
* feat: wire outbox worker in financial-accounting service
Integrate the transactional outbox pattern worker into the financial-accounting
service for reliable event publishing with at-least-once delivery guarantees.
Changes:
- Add ControlAction type and ControlLog domain method for SUSPEND/RESUME/TERMINATE
lifecycle operations following BIAN CoCR patterns
- Add WithTransaction and DB methods to LedgerRepository for atomic outbox writes
- Add protobuf definitions for ControlFinancialBookingLog RPC and events
- Wire OutboxRepository and Worker in main.go with Kafka producer
- Implement ControlFinancialBookingLog gRPC method with idempotency support
- Add comprehensive tests for control operations and state machine logic
The outbox worker polls for pending events and publishes to Kafka, with
configurable batch size and poll interval. Events are written atomically
with domain state changes in a single transaction.
* fix: address CodeRabbit review feedback
Addresses three critical issues identified by CodeRabbit:
1. Kafka producer flush: Call Flush() with 5s timeout before Close() to
ensure all pending outbox events are delivered before shutdown.
2. Race condition fix: Eliminate double-fetch by performing all operations
within a single transaction with pessimistic locking (SELECT FOR UPDATE).
Previously fetched entity outside transaction for domain logic, then
again inside transaction, creating window for concurrent modifications.
Now fetch-lock-apply-save happens atomically.
3. Layer separation: Service layer reconstructs domain model from locked
entity, applies domain logic, then updates entity - maintaining proper
separation while avoiding race conditions.
The refactored ControlBookingLog method now:
- Acquires pessimistic lock on entity within transaction
- Reconstructs domain model from locked entity
- Applies domain control logic
- Updates locked entity with domain results
- Writes event to outbox atomically
- Maps all domain errors to gRPC codes after transaction
---------
Co-authored-by: Ben Coombs <bjcoombs@users.noreply.github.com>1 parent 6379be9 commit fa860fd
12 files changed
Lines changed: 1201 additions & 61 deletions
File tree
- api/proto/meridian
- events/v1
- financial_accounting/v1
- services/financial-accounting
- adapters/persistence
- cmd
- domain
- service
Lines changed: 59 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
452 | 452 | | |
453 | 453 | | |
454 | 454 | | |
| 455 | + | |
| 456 | + | |
| 457 | + | |
| 458 | + | |
| 459 | + | |
| 460 | + | |
| 461 | + | |
| 462 | + | |
| 463 | + | |
| 464 | + | |
| 465 | + | |
| 466 | + | |
| 467 | + | |
| 468 | + | |
| 469 | + | |
| 470 | + | |
| 471 | + | |
| 472 | + | |
| 473 | + | |
| 474 | + | |
| 475 | + | |
| 476 | + | |
| 477 | + | |
| 478 | + | |
| 479 | + | |
| 480 | + | |
| 481 | + | |
| 482 | + | |
| 483 | + | |
| 484 | + | |
| 485 | + | |
| 486 | + | |
| 487 | + | |
| 488 | + | |
| 489 | + | |
| 490 | + | |
| 491 | + | |
| 492 | + | |
| 493 | + | |
| 494 | + | |
| 495 | + | |
| 496 | + | |
| 497 | + | |
| 498 | + | |
| 499 | + | |
| 500 | + | |
| 501 | + | |
| 502 | + | |
| 503 | + | |
| 504 | + | |
| 505 | + | |
| 506 | + | |
| 507 | + | |
| 508 | + | |
| 509 | + | |
| 510 | + | |
| 511 | + | |
| 512 | + | |
| 513 | + | |
455 | 514 | | |
456 | 515 | | |
457 | 516 | | |
| |||
Lines changed: 61 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
104 | 104 | | |
105 | 105 | | |
106 | 106 | | |
| 107 | + | |
| 108 | + | |
| 109 | + | |
| 110 | + | |
| 111 | + | |
| 112 | + | |
| 113 | + | |
| 114 | + | |
| 115 | + | |
| 116 | + | |
| 117 | + | |
| 118 | + | |
| 119 | + | |
107 | 120 | | |
108 | 121 | | |
109 | 122 | | |
| |||
429 | 442 | | |
430 | 443 | | |
431 | 444 | | |
| 445 | + | |
| 446 | + | |
| 447 | + | |
| 448 | + | |
| 449 | + | |
| 450 | + | |
| 451 | + | |
| 452 | + | |
| 453 | + | |
| 454 | + | |
| 455 | + | |
| 456 | + | |
| 457 | + | |
| 458 | + | |
| 459 | + | |
| 460 | + | |
| 461 | + | |
| 462 | + | |
| 463 | + | |
| 464 | + | |
| 465 | + | |
| 466 | + | |
| 467 | + | |
| 468 | + | |
| 469 | + | |
| 470 | + | |
| 471 | + | |
| 472 | + | |
| 473 | + | |
| 474 | + | |
| 475 | + | |
| 476 | + | |
| 477 | + | |
| 478 | + | |
| 479 | + | |
| 480 | + | |
| 481 | + | |
| 482 | + | |
432 | 483 | | |
433 | 484 | | |
434 | 485 | | |
| |||
506 | 557 | | |
507 | 558 | | |
508 | 559 | | |
| 560 | + | |
| 561 | + | |
| 562 | + | |
| 563 | + | |
| 564 | + | |
| 565 | + | |
| 566 | + | |
| 567 | + | |
| 568 | + | |
| 569 | + | |
509 | 570 | | |
Lines changed: 26 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
666 | 666 | | |
667 | 667 | | |
668 | 668 | | |
| 669 | + | |
| 670 | + | |
| 671 | + | |
| 672 | + | |
| 673 | + | |
| 674 | + | |
| 675 | + | |
| 676 | + | |
| 677 | + | |
| 678 | + | |
| 679 | + | |
| 680 | + | |
| 681 | + | |
| 682 | + | |
| 683 | + | |
| 684 | + | |
| 685 | + | |
| 686 | + | |
| 687 | + | |
| 688 | + | |
| 689 | + | |
| 690 | + | |
| 691 | + | |
| 692 | + | |
| 693 | + | |
| 694 | + | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
14 | 14 | | |
15 | 15 | | |
16 | 16 | | |
| 17 | + | |
17 | 18 | | |
18 | 19 | | |
19 | 20 | | |
| |||
22 | 23 | | |
23 | 24 | | |
24 | 25 | | |
| 26 | + | |
25 | 27 | | |
26 | 28 | | |
27 | 29 | | |
| |||
39 | 41 | | |
40 | 42 | | |
41 | 43 | | |
| 44 | + | |
| 45 | + | |
| 46 | + | |
| 47 | + | |
| 48 | + | |
42 | 49 | | |
43 | 50 | | |
44 | 51 | | |
| |||
126 | 133 | | |
127 | 134 | | |
128 | 135 | | |
| 136 | + | |
| 137 | + | |
| 138 | + | |
| 139 | + | |
| 140 | + | |
| 141 | + | |
| 142 | + | |
| 143 | + | |
| 144 | + | |
| 145 | + | |
| 146 | + | |
| 147 | + | |
| 148 | + | |
| 149 | + | |
| 150 | + | |
| 151 | + | |
| 152 | + | |
| 153 | + | |
| 154 | + | |
| 155 | + | |
| 156 | + | |
| 157 | + | |
| 158 | + | |
| 159 | + | |
| 160 | + | |
| 161 | + | |
| 162 | + | |
| 163 | + | |
| 164 | + | |
| 165 | + | |
| 166 | + | |
| 167 | + | |
| 168 | + | |
| 169 | + | |
| 170 | + | |
| 171 | + | |
| 172 | + | |
| 173 | + | |
| 174 | + | |
| 175 | + | |
| 176 | + | |
129 | 177 | | |
130 | 178 | | |
131 | 179 | | |
| |||
171 | 219 | | |
172 | 220 | | |
173 | 221 | | |
174 | | - | |
| 222 | + | |
| 223 | + | |
| 224 | + | |
| 225 | + | |
| 226 | + | |
| 227 | + | |
| 228 | + | |
175 | 229 | | |
176 | 230 | | |
177 | 231 | | |
| |||
314 | 368 | | |
315 | 369 | | |
316 | 370 | | |
317 | | - | |
| 371 | + | |
| 372 | + | |
| 373 | + | |
| 374 | + | |
| 375 | + | |
| 376 | + | |
| 377 | + | |
| 378 | + | |
| 379 | + | |
| 380 | + | |
| 381 | + | |
| 382 | + | |
| 383 | + | |
| 384 | + | |
| 385 | + | |
| 386 | + | |
| 387 | + | |
| 388 | + | |
318 | 389 | | |
319 | 390 | | |
320 | 391 | | |
| |||
0 commit comments