107 lines
3.1 KiB
Go
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`
|