4 Commits

Author SHA256 Message Date
nfel 01d13dbbf2 feat: update settlement workflow
CI/CD Stage / stage (push) Failing after 37s
2026-09-22 16:26:16 +03:30
nfel 4ae5bb615b feat: add Kuknos submission observability
CI/CD Stage / stage (push) Failing after 36s
2026-09-22 16:01:23 +03:30
nfel 837fe5bc5f fix(ci): support SHA-256 repository checkout
CI/CD Stage / stage (push) Failing after 1m20s
2026-09-22 12:41:12 +03:30
nfel 495a5884df Merge feat/refactor-v1 into stage
CI/CD Stage / stage (push) Failing after 31s
2026-09-22 11:45:49 +03:30
6 changed files with 113 additions and 14 deletions
+1 -1
View File
@@ -9,7 +9,7 @@ jobs:
runs-on: ubuntu-latest
steps:
- name: Checkout Git
uses: https://git.darano.ir/actions/checkout@v5
uses: https://git.darano.ir/actions/checkout@v6.0.3
with:
token: ${{ gitea.token }}
path: ./
+1 -1
View File
@@ -9,7 +9,7 @@ jobs:
runs-on: ubuntu-latest
steps:
- name: Checkout Git
uses: https://git.darano.ir/actions/checkout@v5
uses: https://git.darano.ir/actions/checkout@v6.0.3
with:
token: ${{ gitea.token }}
path: ./
+1 -1
View File
@@ -9,7 +9,7 @@ jobs:
runs-on: ubuntu-latest
steps:
- name: Checkout Git
uses: https://git.darano.ir/actions/checkout@v5
uses: https://git.darano.ir/actions/checkout@v6.0.3
with:
token: ${{ gitea.token }}
path: ./
+8 -1
View File
@@ -2,6 +2,7 @@ package settlement
import (
"context"
"errors"
"fmt"
"time"
@@ -41,7 +42,13 @@ func (w Worker) RunOnce(ctx context.Context, limit int) error {
record.Attempts++
hash, submitErr := w.Submitter.Submit(ctx, record)
if submitErr != nil {
record = record.Retry(now, w.MaxAttempts, submitErr)
var permanent ledger.PermanentSubmissionError
if errors.As(submitErr, &permanent) {
record.Status = ledger.SettlementManualReview
record.LastError = submitErr.Error()
} else {
record = record.Retry(now, w.MaxAttempts, submitErr)
}
} else {
// Horizon returns success only after the transaction is accepted into
// a ledger, so a successful response is already confirmation.
+5
View File
@@ -5,6 +5,11 @@ import (
"time"
)
type PermanentSubmissionError struct{ Err error }
func (e PermanentSubmissionError) Error() string { return e.Err.Error() }
func (e PermanentSubmissionError) Unwrap() error { return e.Err }
type SettlementStatus string
const (
+97 -10
View File
@@ -4,13 +4,22 @@ import (
"context"
"encoding/json"
"fmt"
"io"
"log/slog"
"net/http"
"net/url"
"strings"
"time"
"gl/domain/ledger"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/codes"
"go.opentelemetry.io/otel/trace"
)
var tracer = otel.Tracer("gl/infrastructure/kuknos")
// Submitter sends a pre-signed Stellar/Kuknos transaction to Horizon. Signing
// remains outside GL; GL owns durable admission, retry, and reconciliation.
type Submitter struct {
@@ -19,25 +28,51 @@ type Submitter struct {
}
func (s Submitter) Submit(ctx context.Context, record ledger.Settlement) (string, error) {
ctx, span := tracer.Start(ctx, "kuknos.submit_transaction", trace.WithSpanKind(trace.SpanKindClient))
defer span.End()
span.SetAttributes(
attribute.String("kuknos.network", record.Network),
attribute.String("settlement.id", record.ID),
attribute.String("settlement.source_service", record.SourceService),
attribute.String("settlement.source_transaction_id", record.SourceTxID),
attribute.String("settlement.idempotency_key", record.IdempotencyKey),
attribute.Int("settlement.attempt", record.Attempts),
)
if s.Endpoint == "" || record.SignedTransactionXDR == "" {
return "", fmt.Errorf("kuknos submission endpoint and signed xdr are required")
return s.fail(ctx, span, record, fmt.Errorf("kuknos submission endpoint and signed xdr are required"))
}
endpoint := strings.TrimRight(s.Endpoint, "/") + "/transactions"
form := url.Values{"tx": {record.SignedTransactionXDR}}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, strings.TrimRight(s.Endpoint, "/")+"/transactions", strings.NewReader(form.Encode()))
req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, strings.NewReader(form.Encode()))
if err != nil {
return "", err
return s.fail(ctx, span, record, err)
}
req.Header.Set("Content-Type", "application/x-www-form-urlencoded")
client := s.Client
if client == nil {
client = http.DefaultClient
}
startedAt := time.Now()
slog.InfoContext(ctx, "kuknos transaction submission started",
"settlement_id", record.ID,
"source_transaction_id", record.SourceTxID,
"idempotency_key", record.IdempotencyKey,
"network", record.Network,
"attempt", record.Attempts,
"endpoint", endpoint,
)
resp, err := client.Do(req)
if err != nil {
return "", err
return s.fail(ctx, span, record, err)
}
defer resp.Body.Close()
var body struct {
responseBody, err := io.ReadAll(resp.Body)
if err != nil {
return s.fail(ctx, span, record, fmt.Errorf("read kuknos response: %w", err))
}
rawResponse := string(responseBody)
var response struct {
Hash string `json:"hash"`
Extras struct {
ResultCodes struct {
@@ -45,11 +80,63 @@ func (s Submitter) Submit(ctx context.Context, record ledger.Settlement) (string
} `json:"result_codes"`
} `json:"extras"`
}
if err := json.NewDecoder(resp.Body).Decode(&body); err != nil {
return "", fmt.Errorf("decode kuknos response: %w", err)
if err := json.Unmarshal(responseBody, &response); err != nil {
slog.ErrorContext(ctx, "kuknos transaction submission returned an invalid response",
"settlement_id", record.ID,
"source_transaction_id", record.SourceTxID,
"idempotency_key", record.IdempotencyKey,
"network", record.Network,
"attempt", record.Attempts,
"status_code", resp.StatusCode,
"stellar_response", rawResponse,
"error", err,
)
return s.fail(ctx, span, record, fmt.Errorf("decode kuknos response: %w", err))
}
if resp.StatusCode < 200 || resp.StatusCode >= 300 || body.Hash == "" {
return "", fmt.Errorf("kuknos rejected transaction: status=%d code=%s", resp.StatusCode, body.Extras.ResultCodes.Transaction)
latency := time.Since(startedAt)
span.SetAttributes(
attribute.Int("http.response.status_code", resp.StatusCode),
attribute.Int64("http.response.body.size", int64(len(responseBody))),
attribute.Int64("kuknos.submission.duration_ms", latency.Milliseconds()),
)
span.AddEvent("kuknos.response", trace.WithAttributes(
attribute.String("stellar_response", rawResponse),
))
span.SetStatus(codes.Ok, "submission response received")
slog.InfoContext(ctx, "kuknos transaction submission response",
"settlement_id", record.ID,
"source_transaction_id", record.SourceTxID,
"idempotency_key", record.IdempotencyKey,
"network", record.Network,
"attempt", record.Attempts,
"status_code", resp.StatusCode,
"transaction_hash", response.Hash,
"transaction_result_code", response.Extras.ResultCodes.Transaction,
"stellar_response", rawResponse,
"duration_ms", latency.Milliseconds(),
)
if resp.StatusCode < 200 || resp.StatusCode >= 300 || response.Hash == "" {
err := fmt.Errorf("kuknos rejected transaction: status=%d code=%s", resp.StatusCode, response.Extras.ResultCodes.Transaction)
span.SetStatus(codes.Error, err.Error())
span.RecordError(err)
if response.Extras.ResultCodes.Transaction == "tx_bad_seq" {
return "", ledger.PermanentSubmissionError{Err: err}
}
return "", err
}
return body.Hash, nil
return response.Hash, nil
}
func (s Submitter) fail(ctx context.Context, span trace.Span, record ledger.Settlement, err error) (string, error) {
span.SetStatus(codes.Error, err.Error())
span.RecordError(err)
slog.ErrorContext(ctx, "kuknos transaction submission failed",
"settlement_id", record.ID,
"source_transaction_id", record.SourceTxID,
"idempotency_key", record.IdempotencyKey,
"network", record.Network,
"attempt", record.Attempts,
"error", err,
)
return "", err
}