Files

107 lines
3.1 KiB
Go

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`