rdma: harden RC outcome publication ownership and record fidelity

Review of the publication path found that ownership could change
hands at the wrong moment and that records could disagree with both
the wire response and the underlying operation.

Successful transfers lost their completion record: the native
completion calls fire the teardown callback synchronously, so the
callback claimed the publication first and logged every completed
transfer as an expiry, and committed PUTs produced no object-created
events. The READY handler now reserves the publication before
invoking any completion call; a reserved record is invisible to the
callback, and the handler publishes the real outcome exactly once.

The same race existed at creation: the PREPARE finalizer can reap an
expired session and fire the callback before the session is
registered, leaving an orphan entry whose only notification already
happened. Registration now runs before the finalizing call, a
notification that arrives first is parked and consumed by the
registration, and a failed finalization drops the entry.

Retained records referenced the request's pooled header buffer, so
a later request could rewrite a tracked session's bucket and key;
captured strings are cloned now. The synthesized publication path
follows the same rule for the event senders, which serialize
asynchronously.

Operational sinks classify plain errors as 500 on their own, so a
resource-limit rejection logged 500 while the client saw 429. The
publication renders non-S3 errors through the route error mapping
before the record reaches the sinks, and the expiry record carries
a dedicated SessionExpired code instead of a generic one. Malformed
PUT headers now preserve the operation in the record.
This commit is contained in:
Jihyeon Gim
2026-09-09 13:16:52 +09:00
parent 39aae5435c
commit 279f5e7d2f
3 changed files with 282 additions and 74 deletions
+175 -57
View File
@@ -18,7 +18,8 @@
package rcroutes
import (
"fmt"
"errors"
"strings"
"sync"
"time"
@@ -29,6 +30,7 @@ import (
"github.com/versity/versitygw/metrics"
"github.com/versity/versitygw/rdma/rcserver"
"github.com/versity/versitygw/s3api/utils"
"github.com/versity/versitygw/s3err"
"github.com/versity/versitygw/s3event"
"github.com/versity/versitygw/s3log"
)
@@ -45,6 +47,9 @@ type OpsServices struct {
// opsEmitter is the operational context captured at PREPARE and held
// until the final outcome is known: enough to synthesize an access
// record carrying the session's object rather than the wire path.
// Every string field is owned storage: nothing may reference the
// request's pooled buffers once PREPARE returns, because fasthttp
// reuses them for the next request.
type opsEmitter struct {
ops OpsServices
app *fiber.App
@@ -59,7 +64,9 @@ type opsEmitter struct {
// synthesize builds a fiber context whose path and request locals
// describe the session's logical object operation, so the standard
// access-log and event pipelines observe GET/PUT of bucket/key
// instead of the fixed RDMA control path.
// instead of the fixed RDMA control path. The path string is the
// emitter's own storage: event senders serialize asynchronously, so
// the synthesized context must never hand them pooled buffers.
func (e *opsEmitter) synthesize() (fiber.Ctx, func()) {
ctx := e.app.AcquireCtx(&fasthttp.RequestCtx{})
method := fiber.MethodGet
@@ -69,7 +76,9 @@ func (e *opsEmitter) synthesize() (fiber.Ctx, func()) {
ctx.Method(method)
// The access logger and the event schema both split this path
// into bucket/key, so the synthesized path must be the object
// path in canonical form.
// path in canonical form. fiber copies override strings it
// stores as the path original; the derived c.path below is
// a fresh allocation, which is what outlives the release.
ctx.Path("/" + e.bucket + "/" + e.key)
utils.ContextKeyAccount.Set(ctx, e.acct)
utils.ContextKeyRegion.Set(ctx, e.region)
@@ -84,10 +93,23 @@ func (e *opsEmitter) synthesize() (fiber.Ctx, func()) {
// publish emits the final audit record, request metric, and (for a
// committed PUT) the object-created event. Exactly-once delivery is
// the tracker's job; this method just performs one emission.
//
// The operational sinks classify plain errors as 500 on their own,
// which would disagree with the status the client saw. Before the
// record reaches them, the error is rendered as its mapped S3 error,
// so the audit log, the metric, and the wire response all carry the
// same classification.
func (e *opsEmitter) publish(err error, bytes int64) {
if e == nil || (e.ops.Logger == nil && e.ops.Metrics == nil && e.ops.Events == nil) {
return
}
sinkErr := err
if err != nil {
var s3Err s3err.S3Error
if !errors.As(err, &s3Err) {
sinkErr = routeError(err)
}
}
ctx, release := e.synthesize()
defer release()
@@ -95,13 +117,13 @@ func (e *opsEmitter) publish(err error, bytes int64) {
if e.isPut {
action = metrics.ActionPutObject
}
status := httpStatusFromError(err)
status := httpStatusFromError(sinkErr)
if e.ops.Metrics != nil {
e.ops.Metrics.Send(ctx, err, action, bytes, status)
e.ops.Metrics.Send(ctx, sinkErr, action, bytes, status)
}
if e.ops.Logger != nil {
e.ops.Logger.Log(ctx, err, nil, s3log.LogMeta{
e.ops.Logger.Log(ctx, sinkErr, nil, s3log.LogMeta{
Action: action,
})
}
@@ -116,9 +138,8 @@ func (e *opsEmitter) publish(err error, bytes int64) {
}
// httpStatusFromError maps an operation error to the HTTP status
// the S3 surface would have answered with, using the same route
// error mapping as the wire response so operational records never
// disagree with what the client saw.
// the S3 surface would have answered with. The callers pass mapped
// S3 errors, so the status is simply the error's own.
func httpStatusFromError(err error) int {
if err == nil {
return 200
@@ -128,81 +149,157 @@ func httpStatusFromError(err error) int {
// sessionRecord is one tracked session with its captured context.
type sessionRecord struct {
emit *opsEmitter
done bool
pending bool
emit *opsEmitter
// reserved marks a record the request path holds exclusively:
// it took ownership before invoking a native completion call
// that would fire the teardown callback synchronously, so the
// callback must not publish on its behalf.
reserved bool
}
// opsTracker owns terminal publication: each session (and each
// pre-session request) is published exactly once, by whichever path
// confirms the final outcome first. It carries its own throwaway
// fiber app: the synthesized contexts only carry path and locals,
// never route state, so they must not share the gateway app.
// fiber.App for synthesizing publication contexts, independent of
// the gateway's request routing.
type opsTracker struct {
mu sync.Mutex
sessions map[string]*sessionRecord
ops OpsServices
app *fiber.App
sessions map[string]*sessionRecord
// earlyTerminals parks teardown notifications that arrived
// before the session's registration; register consumes them.
earlyTerminals map[string]rcserver.TerminalEvent
app *fiber.App
}
func newOpsTracker() *opsTracker {
return &opsTracker{
sessions: map[string]*sessionRecord{},
app: fiber.New(),
sessions: map[string]*sessionRecord{},
earlyTerminals: map[string]rcserver.TerminalEvent{},
app: fiber.New(),
}
}
// SetOpsServices installs the operational service instances. The
// gateway creates the logger, metrics manager, and event sender
// after the RC routes exist, so the tracker starts empty and the
// embedder injects them once RunVersityGW has built them. Sessions
// registered before the injection publish nothing (there are none:
// the server is not listening yet).
// services arrive here. Sessions registered before the injection
// publish nothing (there are none: the gateway wires this before
// it starts serving).
func (t *opsTracker) SetOpsServices(ops OpsServices) {
t.mu.Lock()
defer t.mu.Unlock()
t.ops = ops
}
// opsSnapshot returns the current service set under the lock.
func (t *opsTracker) opsSnapshot() OpsServices {
// register captures the operational context of a successfully created
// session so a later terminal path can publish its final outcome.
// The strings are cloned: they originate from the request's pooled
// header buffer, which does not survive the response.
//
// handleEarlyTerminal covers the FinishPrepare race: the native call
// that finalizes PREPARE can reap an already-expired session and fire
// the teardown callback before register runs. When the callback wins
// that race it parks the event, and register consumes it instead of
// leaving an entry whose only notification already happened.
func (t *opsTracker) register(sessionID string, acct auth.Account,
region, bucket, key string, isPut bool, start time.Time) {
emit := &opsEmitter{
ops: t.loadOps(),
app: t.app,
acct: acct,
region: region,
bucket: strings.Clone(bucket),
key: strings.Clone(key),
isPut: isPut,
start: start,
}
t.mu.Lock()
defer t.mu.Unlock()
// A teardown notification that arrived before this registration
// owns the publication: publish now and store nothing.
if _, parked := t.earlyTerminals[sessionID]; parked {
delete(t.earlyTerminals, sessionID)
go emit.publish(errSessionExpired, 0)
return
}
t.sessions[sessionID] = &sessionRecord{emit: emit}
}
// unregister drops a session entry whose PREPARE finalization
// failed: the native side is gone, so the teardown callback has
// either already published or will find nothing. A parked early
// notification is dropped with it (the failure publication covers
// the outcome).
func (t *opsTracker) unregister(sessionID string) {
if t == nil {
return
}
t.mu.Lock()
defer t.mu.Unlock()
delete(t.sessions, sessionID)
delete(t.earlyTerminals, sessionID)
}
// reserve takes exclusive ownership of a session's publication
// before the request path invokes a native completion call
// (FinishFinal, FinishPut, or a reap-triggering mutation). Those
// calls fire the teardown callback synchronously while the session
// record is still live; reserving first keeps the callback from
// publishing an expiry record for a transfer that is completing
// right now.
func (t *opsTracker) reserve(sessionID string) *opsEmitter {
t.mu.Lock()
defer t.mu.Unlock()
rec, ok := t.sessions[sessionID]
if !ok || rec.reserved {
return nil
}
rec.reserved = true
return rec.emit
}
// unreserve restores callback ownership when a reserved completion
// call did not after all retire the session (the caller failed
// before any state change). The record goes back to normal tracking
// unless a teardown notification landed meanwhile.
func (t *opsTracker) unreserve(sessionID string, emit *opsEmitter) {
if t == nil || emit == nil {
return
}
t.mu.Lock()
defer t.mu.Unlock()
rec, ok := t.sessions[sessionID]
if !ok {
// The session is gone: the completion call retired it and
// the reserved emitter is the only remaining owner, so
// nothing to restore.
return
}
rec.reserved = false
}
func (t *opsTracker) loadOps() OpsServices {
t.mu.Lock()
defer t.mu.Unlock()
return t.ops
}
// register captures the operational context of a successfully created
// session so a later terminal path can publish its final outcome.
func (t *opsTracker) register(sessionID string, acct auth.Account,
region, bucket, key string, isPut bool, start time.Time) {
emit := &opsEmitter{
ops: t.opsSnapshot(),
app: t.app,
acct: acct,
region: region,
bucket: bucket,
key: key,
isPut: isPut,
start: start,
}
t.mu.Lock()
defer t.mu.Unlock()
t.sessions[sessionID] = &sessionRecord{emit: emit, pending: true}
}
// claim removes the session's publication slot and returns its
// captured context; the second caller gets nil and publishes nothing.
func (t *opsTracker) claim(sessionID string) *opsEmitter {
// claim removes the session from the table and returns its emitter
// to exactly one publisher. A reserved record is only claimable by
// its reserving request path (the callback skips it).
func (t *opsTracker) claim(sessionID string, byRequest bool) *opsEmitter {
t.mu.Lock()
defer t.mu.Unlock()
rec, ok := t.sessions[sessionID]
if !ok {
return nil
}
delete(t.sessions, sessionID)
if !rec.pending {
if rec.reserved && !byRequest {
return nil
}
delete(t.sessions, sessionID)
return rec.emit
}
@@ -213,13 +310,22 @@ func (t *opsTracker) onTerminal(ev rcserver.TerminalEvent) {
if t == nil {
return
}
emit := t.claim(ev.SessionID)
if emit == nil {
emit := t.claim(ev.SessionID, false)
if emit != nil {
// An expired or abandoned session never reached a final
// object result; the bytes staged for it did not become
// a transfer.
emit.publish(errSessionExpired, 0)
return
}
// An expired or abandoned session never reached a final object
// result; the bytes staged for it did not become a transfer.
emit.publish(errSessionExpired, 0)
// A session tearing down before its registration ran: park the
// event so register can publish instead of orphaning a record
// whose only notification already happened.
t.mu.Lock()
if _, live := t.sessions[ev.SessionID]; !live {
t.earlyTerminals[ev.SessionID] = ev
}
t.mu.Unlock()
}
// publishClaimed publishes the terminal record from the request
@@ -230,7 +336,7 @@ func (t *opsTracker) publishClaimed(sessionID string, err error, bytes ...int64)
if t == nil {
return
}
emit := t.claim(sessionID)
emit := t.claim(sessionID, true)
if emit == nil {
return
}
@@ -250,12 +356,12 @@ func (t *opsTracker) publishRequest(ctx fiber.Ctx, acct auth.Account,
return
}
emit := &opsEmitter{
ops: t.ops,
ops: t.loadOps(),
app: t.app,
acct: acct,
region: regionFromCtx(ctx),
bucket: bucket,
key: key,
bucket: strings.Clone(bucket),
key: strings.Clone(key),
isPut: isPut,
start: time.Now(),
}
@@ -271,4 +377,16 @@ func regionFromCtx(ctx fiber.Ctx) string {
return ""
}
var errSessionExpired = fmt.Errorf("session expired")
// expiredOutcomeError renders a parked teardown notification as the
// error the publication carries: an internal S3 error whose code
// names the expiry, so the audit log keeps the descriptive code the
// plain error used to carry instead of the generic mapping.
type sessionExpiredError struct {
s3err.APIError
}
var errSessionExpired = sessionExpiredError{APIError: s3err.APIError{
Code: "SessionExpired",
Description: "The RDMA transfer session expired before completion",
HTTPStatusCode: 500,
}}
+70
View File
@@ -78,6 +78,71 @@ func TestOpsTrackerRequestPathClaims(t *testing.T) {
}
}
func TestOpsTrackerReserveBlocksCallback(t *testing.T) {
tr := newOpsTracker()
tr.register("sess-3", auth.Account{Access: "ak"}, "us-east-1",
"bkt", "obj", true, time.Now())
// The request path reserves before invoking a native
// completion call.
emit := tr.reserve("sess-3")
if emit == nil {
t.Fatal("reserve returned nil for a live session")
}
// The synchronous teardown callback must not publish on a
// reserved record: the entry stays put.
tr.onTerminal(rcserver.TerminalEvent{SessionID: "sess-3"})
if got := len(tr.sessions); got != 1 {
t.Fatalf("reserved session removed by callback: %d", got)
}
// A double reserve is refused.
if tr.reserve("sess-3") != nil {
t.Fatal("double reserve succeeded")
}
// The request path publishes through its reserved emitter.
tr.publishClaimed("sess-3", nil, 128)
if got := len(tr.sessions); got != 0 {
t.Fatalf("reserved session survived request publication: %d", got)
}
}
func TestOpsTrackerEarlyTerminal(t *testing.T) {
tr := newOpsTracker()
// Teardown notification for an unregistered session parks.
tr.onTerminal(rcserver.TerminalEvent{SessionID: "sess-4"})
if got := len(tr.earlyTerminals); got != 1 {
t.Fatalf("early terminal not parked: %d", got)
}
// Registration consumes it: no live entry remains, so no
// orphan record can outlive the session.
tr.register("sess-4", auth.Account{Access: "ak"}, "us-east-1",
"bkt", "obj", false, time.Now())
if got := len(tr.earlyTerminals); got != 0 {
t.Fatalf("early terminal not consumed: %d", got)
}
if got := len(tr.sessions); got != 0 {
t.Fatalf("orphan session entry created: %d", got)
}
}
func TestOpsTrackerUnregister(t *testing.T) {
tr := newOpsTracker()
tr.register("sess-5", auth.Account{Access: "ak"}, "us-east-1",
"bkt", "obj", false, time.Now())
tr.onTerminal(rcserver.TerminalEvent{SessionID: "sess-5"})
tr.register("sess-6", auth.Account{Access: "ak"}, "us-east-1",
"bkt", "obj", false, time.Now())
tr.unregister("sess-6")
if got := len(tr.sessions); got != 0 {
t.Fatalf("unregister left entries: %d", got)
}
}
func TestOpsTrackerUnknownSession(t *testing.T) {
tr := newOpsTracker()
// Unknown sessions and the nil tracker are silent no-ops.
@@ -101,4 +166,9 @@ func TestHttpStatusFromError(t *testing.T) {
if got := httpStatusFromError(s3err.GetAPIError(s3err.ErrAccessDenied)); got != 403 {
t.Fatalf("access denied => %d, want 403", got)
}
// A resource-limit rejection maps to the wire status, not a
// generic 500.
if got := httpStatusFromError(rcserver.ErrLimit); got != 429 {
t.Fatalf("limit error => %d, want 429", got)
}
}
+37 -17
View File
@@ -155,43 +155,43 @@ func (h *Handler) prepareCore(ctx fiber.Ctx) error {
// Header parse failures end the request before authorization;
// publish them as request records too, with whatever object
// identity the malformed headers still carried.
publishHeaderErr := func(err error) error {
// identity and operation the malformed headers still carried.
publishHeaderErr := func(err error, isPut bool) error {
target := ctx.Get(hdrTarget)
bucket, key, _ := splitTarget(target)
h.ops.publishRequest(ctx, acct, err, bucket, key, false)
h.ops.publishRequest(ctx, acct, err, bucket, key, isPut)
return err
}
if proto := ctx.Get(hdrProtocol); proto != protocolValue {
return publishHeaderErr(invalidHeader(hdrProtocol, proto))
return publishHeaderErr(invalidHeader(hdrProtocol, proto), false)
}
op := strings.ToUpper(ctx.Get(hdrOp))
if op != "GET" && op != "PUT" {
return publishHeaderErr(invalidHeader(hdrOp, ctx.Get(hdrOp)))
return publishHeaderErr(invalidHeader(hdrOp, ctx.Get(hdrOp)), false)
}
isPut := op == "PUT"
target := ctx.Get(hdrTarget)
bucket, key, ok := splitTarget(target)
if !ok {
return publishHeaderErr(invalidHeader(hdrTarget, target))
return publishHeaderErr(invalidHeader(hdrTarget, target), isPut)
}
size, err := parseUint(ctx.Get(hdrSize), 10, 64)
if err != nil || size == 0 {
return publishHeaderErr(invalidHeader(hdrSize, ctx.Get(hdrSize)))
return publishHeaderErr(invalidHeader(hdrSize, ctx.Get(hdrSize)), isPut)
}
offset, err := parseUint(ctx.Get(hdrOffset), 10, 64)
if err != nil {
return publishHeaderErr(invalidHeader(hdrOffset, ctx.Get(hdrOffset)))
return publishHeaderErr(invalidHeader(hdrOffset, ctx.Get(hdrOffset)), isPut)
}
psn, err := parseUint(ctx.Get(hdrPsn), 16, 32)
if err != nil || psn == 0 || psn > 0xffffff {
return publishHeaderErr(invalidHeader(hdrPsn, ctx.Get(hdrPsn)))
return publishHeaderErr(invalidHeader(hdrPsn, ctx.Get(hdrPsn)), isPut)
}
cookie, err := parseUint(ctx.Get(hdrCookie), 16, 32)
if err != nil || cookie == 0 {
return publishHeaderErr(invalidHeader(hdrCookie, ctx.Get(hdrCookie)))
return publishHeaderErr(invalidHeader(hdrCookie, ctx.Get(hdrCookie)), isPut)
}
// Authorize through the regular object-access chain.
@@ -224,7 +224,18 @@ func (h *Handler) prepareCore(ctx fiber.Ctx) error {
return err
}
}
// Register before the finalizing call: FinishPrepare can reap
// an already-expired session and fire the teardown callback
// synchronously, and a registered record (or a parked early
// notification) keeps that publication from being lost.
h.ops.register(resp.SessionID, acct,
regionFromCtx(ctx), bucket, key, isPut, time.Now())
if err := h.svc.FinishPrepare(resp.SessionID, true); err != nil {
// The finalization failed: the session is gone and the
// teardown callback (or this call's own reap) published
// the terminal record; nothing may keep the entry.
h.ops.unregister(resp.SessionID)
h.ops.publishRequest(ctx, acct, mapRcError(err), bucket, key, isPut)
return mapRcError(err)
}
@@ -232,10 +243,6 @@ func (h *Handler) prepareCore(ctx fiber.Ctx) error {
// The session now owns the operation record: the terminal
// path (READY/FinishPut completion, CANCEL, or the expiry
// reaper) publishes the final outcome exactly once.
if h.ops != nil {
h.ops.register(resp.SessionID, acct,
regionFromCtx(ctx), bucket, key, isPut, time.Now())
}
// Wire reply per the hipobj-rc-v2 contract: protocol echo,
// the server endpoint as "200:<token>", session id, and PSN.
@@ -409,6 +416,19 @@ func (h *Handler) readyCore(ctx fiber.Ctx) error {
// The claim succeeded: from here until the response commits,
// this handler owns the completion ref. A panic or early
// unwind must still release it so the session can be reaped.
//
// The publication ownership moves here as well: every native
// completion call below (FinishFinal, FinishPut) fires the
// teardown callback synchronously, and a reserved record is
// invisible to that callback, so the outcome is published by
// this handler exactly once.
emit := h.ops.reserve(sessionID)
publish := func(err error, bytes int64) {
if emit != nil {
emit.publish(err, bytes)
emit = nil
}
}
finalized := false
defer func() {
if !finalized {
@@ -429,7 +449,7 @@ func (h *Handler) readyCore(ctx fiber.Ctx) error {
finalized = true
}
if err != nil {
h.ops.publishClaimed(sessionID, err)
publish(mapRcError(err), 0)
return err
}
// The FINAL wire reply carries the stored object's
@@ -437,7 +457,7 @@ func (h *Handler) readyCore(ctx fiber.Ctx) error {
resp.Etag = put.ETag
resp.VersionID = put.VersionID
} else if err := h.svc.FinishFinal(sessionID); err != nil {
h.ops.publishClaimed(sessionID, mapRcError(err))
publish(mapRcError(err), 0)
return mapRcError(err)
} else {
finalized = true
@@ -445,7 +465,7 @@ func (h *Handler) readyCore(ctx fiber.Ctx) error {
// The transfer completed: publish the terminal record with
// the byte count the data plane reported.
h.ops.publishClaimed(sessionID, nil, int64(resp.BytesTransferred))
publish(nil, int64(resp.BytesTransferred))
// Wire reply per the hipobj-rc-v2 contract: protocol echo,
// cookie echo, transferred bytes, and object metadata.