package postgres import ( "context" "encoding/json" "errors" "fmt" "gl/domain/ledger" "github.com/jackc/pgx/v5" ) func (r *JournalRepository) AppendEvent(ctx context.Context, event ledger.TransactionEvent) (ledger.TransactionEvent, bool, error) { if err := event.Validate(); err != nil { return ledger.TransactionEvent{}, false, fmt.Errorf("validate transaction event: %w", err) } metadataValues := event.Metadata if metadataValues == nil { metadataValues = map[string]string{} } metadata, err := json.Marshal(metadataValues) if err != nil { return ledger.TransactionEvent{}, false, fmt.Errorf("encode event metadata: %w", err) } err = r.database.QueryRow(ctx, insertEventSQL, event.ID, event.SourceService, event.IdempotencyKey, event.SourceTransactionID, event.TrackingCode, event.EventVersion, event.State, event.ErrorCode, event.ErrorMessage, event.OccurredAt, event.CorrelationID, event.ActorID, event.Blockchain.Network, event.Blockchain.TransactionHash, event.Blockchain.LedgerSequence, metadata, event.PayloadHash, ).Scan(&event.RecordedAt) if !errors.Is(err, pgx.ErrNoRows) { if err != nil { return ledger.TransactionEvent{}, false, fmt.Errorf("insert transaction event: %w", err) } return event, false, nil } var stored ledger.TransactionEvent var storedMetadata []byte err = r.database.QueryRow(ctx, selectEventByIdempotencySQL, event.IdempotencyKey).Scan( &stored.ID, &stored.SourceService, &stored.IdempotencyKey, &stored.SourceTransactionID, &stored.TrackingCode, &stored.EventVersion, &stored.State, &stored.ErrorCode, &stored.ErrorMessage, &stored.OccurredAt, &stored.RecordedAt, &stored.CorrelationID, &stored.ActorID, &stored.Blockchain.Network, &stored.Blockchain.TransactionHash, &stored.Blockchain.LedgerSequence, &storedMetadata, &stored.PayloadHash, ) if err != nil { return ledger.TransactionEvent{}, false, fmt.Errorf("read idempotent transaction event: %w", err) } if stored.PayloadHash != event.PayloadHash { return ledger.TransactionEvent{}, false, ErrIdempotencyConflict } if err := json.Unmarshal(storedMetadata, &stored.Metadata); err != nil { return ledger.TransactionEvent{}, false, fmt.Errorf("decode event metadata: %w", err) } return stored, true, nil } const insertEventSQL = ` INSERT INTO transaction_events ( id, source_service, idempotency_key, source_transaction_id, tracking_code, event_version, state, error_code, error_message, occurred_at, correlation_id, actor_id, blockchain_network, blockchain_transaction_hash, blockchain_ledger_sequence, metadata, payload_hash ) VALUES ( $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17 ) ON CONFLICT (idempotency_key) DO NOTHING RETURNING recorded_at` const selectEventByIdempotencySQL = ` SELECT id, source_service, idempotency_key, source_transaction_id, tracking_code, event_version, state, error_code, error_message, occurred_at, recorded_at, correlation_id, actor_id, blockchain_network, blockchain_transaction_hash, blockchain_ledger_sequence, metadata, payload_hash FROM transaction_events WHERE idempotency_key = $1`