Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion docs/datasources/dgraph/page.md
Original file line number Diff line number Diff line change
Expand Up @@ -114,7 +114,7 @@ func DGraphInsertHandler(c *gofr.Context) (any, error) {
// Create an api.Mutation object
mutation := &api.Mutation{
SetJson: []byte(mutationData), // Set the JSON payload
CommitNow: true, // Auto-commit the transaction
CommitNow: true, // Optional: commits within the mutation RPC. Mutate commits either way.
}

// Run the mutation in Dgraph
Expand Down
2 changes: 2 additions & 0 deletions pkg/gofr/container/datasources.go
Original file line number Diff line number Diff line change
Expand Up @@ -483,6 +483,8 @@ type Dgraph interface {
QueryWithVars(ctx context.Context, query string, vars map[string]string) (any, error)

// Mutate executes a write operation (mutation) in the Dgraph database and returns the result.
// The write is committed before Mutate returns, whether or not CommitNow is set on the
// mutation. Use NewTxn directly to spread several mutations across one transaction.
// Parameters:
// - ctx: The context for the mutation.
// - mu: The mutation operation, usually of type *api.Mutation.
Expand Down
51 changes: 50 additions & 1 deletion pkg/gofr/datasource/dgraph/dgraph.go
Original file line number Diff line number Diff line change
Expand Up @@ -205,6 +205,8 @@ func (d *Client) QueryWithVars(ctx context.Context, query string, vars map[strin
}

// Mutate executes a write operation (mutation) in the Dgraph database and returns the result.
//
Comment thread
Umang01-hash marked this conversation as resolved.
// The write is committed before this returns, whether or not CommitNow is set on the mutation.
func (d *Client) Mutate(ctx context.Context, mu any) (any, error) {
start := time.Now()

Expand All @@ -217,7 +219,7 @@ func (d *Client) Mutate(ctx context.Context, mu any) (any, error) {
}

// Execute mutation
resp, err := d.client.NewTxn().Mutate(tracedCtx, mutation)
resp, err := d.mutateInTxn(tracedCtx, mutation)
duration := time.Since(start).Microseconds()

// Create and log the mutation details
Expand All @@ -239,6 +241,53 @@ func (d *Client) Mutate(ctx context.Context, mu any) (any, error) {
return resp, nil
}

// mutateInTxn runs the mutation and resolves the transaction it opens.
//
// dgo only finishes the transaction when the caller set CommitNow on the mutation — Txn.Mutate
// copies the field into the api.Request it builds, and Txn.Do marks the transaction finished only
// when that field is set. Mutating without it therefore staged the write into a transaction that
// was then abandoned to the garbage collector: nothing was persisted, no error was returned, and
Comment thread
Umang01-hash marked this conversation as resolved.
// the transaction stayed open on the server until Dgraph timed it out. Committing here makes a
// single Mutate call atomic whichever way the caller wrote it.
func (d *Client) mutateInTxn(ctx context.Context, mutation *api.Mutation) (*api.Response, error) {
txn := d.client.NewTxn()

// Discard delegates to Txn.commitOrAbort, which returns immediately once the transaction is
// finished. Every path out of this function leaves it finished — Txn.Do finishes a CommitNow
// mutation, Txn.Commit below finishes the rest, and Txn.Do discards the transaction itself
// when the mutation fails — so the discard costs no round trip. It is here so that a
// transaction is still released if any of those paths stops finishing it.
//
// The cancellation is dropped for the discard alone: a caller whose context is canceled
// between the commit returning and this deferred call would otherwise have a successful
// write report a discard failure in the logs. The deadline-free context only ever covers
// the abort of a transaction that is already being abandoned.
discardCtx := context.WithoutCancel(ctx)

defer func() {
if err := txn.Discard(discardCtx); err != nil {
d.logger.Error("dgraph mutation transaction discard failed: ", err)
}
}()

resp, err := txn.Mutate(ctx, mutation)
if err != nil {
return nil, err
}

// dgo already committed and marked the transaction finished; Txn.Commit returns ErrFinished
// for a finished transaction.
if mutation.CommitNow {
return resp, nil
}

if err := txn.Commit(ctx); err != nil {
return nil, err
}

return resp, nil
}

// Alter applies schema or other changes to the Dgraph database.
func (d *Client) Alter(ctx context.Context, op any) error {
start := time.Now()
Expand Down
2 changes: 2 additions & 0 deletions pkg/gofr/datasource/dgraph/dgraph_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -155,6 +155,7 @@ func Test_Mutate_Success(t *testing.T) {
mutation := &api.Mutation{CommitNow: true}

mockTxn.EXPECT().Mutate(gomock.Any(), mutation).Return(&api.Response{Json: []byte(`{"result": "mutation success"}`)}, nil)
mockTxn.EXPECT().Discard(gomock.Any()).Return(nil)

mockLogger.EXPECT().Debug(gomock.Any())
mockLogger.EXPECT().Debugf("dgraph mutation succeeded in %dµs", gomock.Any())
Expand Down Expand Up @@ -201,6 +202,7 @@ func Test_Mutate_Error(t *testing.T) {
mutation := &api.Mutation{CommitNow: true}

mockTxn.EXPECT().Mutate(gomock.Any(), mutation).Return(nil, errMutationFailed)
mockTxn.EXPECT().Discard(gomock.Any()).Return(nil)

mockLogger.EXPECT().Debug(gomock.Any())
mockLogger.EXPECT().Error("dgraph mutation failed: ", errMutationFailed)
Expand Down
201 changes: 201 additions & 0 deletions pkg/gofr/datasource/dgraph/mutate_txn_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,201 @@
package dgraph

import (
"context"
"errors"
"testing"

"github.com/dgraph-io/dgo/v210/protos/api"
"github.com/stretchr/testify/require"
"go.opentelemetry.io/otel"
"go.uber.org/mock/gomock"
)

var (
errCommitFailed = errors.New("commit failed")
errDiscardFailed = errors.New("discard failed")
)

// recordingTxn counts what Mutate does to the transaction it opens.
//
// A gomock Txn cannot answer "was Commit called?" in this package: setupDB runs ctrl.Finish()
// when it returns, which marks the controller finished before the test body starts, so the
// t.Cleanup verification gomock.NewController installs short-circuits on Controller.finish's
// already-finished branch, and an unmet expectation is never reported. Counting here asserts on
// what happened rather than on the mock library's bookkeeping.
type recordingTxn struct {
mutateErr error
commitErr error
discardErr error

mutations int
commits int
discards int

// discardCtxErr is whatever the context handed to Discard reported. Mutate
// detaches cancellation for the discard alone, so this stays nil even when
// the caller's context is already done.
discardCtxErr error
}

func (r *recordingTxn) Mutate(_ context.Context, _ *api.Mutation) (*api.Response, error) {
r.mutations++

if r.mutateErr != nil {
return nil, r.mutateErr
}

return &api.Response{Json: []byte(`{}`)}, nil
}

func (r *recordingTxn) Commit(context.Context) error {
r.commits++

return r.commitErr
}

func (r *recordingTxn) Discard(ctx context.Context) error {
r.discards++
r.discardCtxErr = ctx.Err()

return r.discardErr
}

func (r *recordingTxn) BestEffort() Txn { return r }

func (*recordingTxn) Query(context.Context, string) (*api.Response, error) { return nil, nil }

func (*recordingTxn) QueryRDF(context.Context, string) (*api.Response, error) { return nil, nil }

func (*recordingTxn) QueryWithVars(context.Context, string, map[string]string) (*api.Response, error) {
return nil, nil
}

func (*recordingTxn) QueryRDFWithVars(context.Context, string,
map[string]string) (*api.Response, error) {
return nil, nil
}

func (*recordingTxn) Do(context.Context, *api.Request) (*api.Response, error) { return nil, nil }

func setupWithTxn(t *testing.T, txn Txn) *Client {
t.Helper()

ctrl := gomock.NewController(t)

logger := NewMockLogger(ctrl)
logger.EXPECT().Debug(gomock.Any()).AnyTimes()
logger.EXPECT().Debugf(gomock.Any(), gomock.Any()).AnyTimes()
logger.EXPECT().Log(gomock.Any()).AnyTimes()
logger.EXPECT().Logf(gomock.Any(), gomock.Any()).AnyTimes()
logger.EXPECT().Error(gomock.Any(), gomock.Any()).AnyTimes()

metrics := NewMockMetrics(ctrl)
metrics.EXPECT().RecordHistogram(gomock.Any(), gomock.Any(), gomock.Any()).AnyTimes()

client := New(Config{Host: "localhost", Port: "9080"})
client.UseLogger(logger)
client.UseMetrics(metrics)
client.UseTracer(otel.GetTracerProvider().Tracer("gofr-dgraph"))

dgraphClient := NewMockDgraphClient(ctrl)
dgraphClient.EXPECT().NewTxn().Return(txn).AnyTimes()
client.client = dgraphClient

return client
}

func Test_Mutate_TransactionHandling(t *testing.T) {
tests := []struct {
name string
commitNow bool
txn recordingTxn
wantErr error
wantResp bool
wantCommits int
wantDiscards int
}{
{
// The defect: dgo only finishes the transaction when CommitNow is set, so without
// an explicit Commit the write was staged and abandoned, silently.
name: "commits when CommitNow is not set",
wantResp: true,
wantCommits: 1,
wantDiscards: 1,
},
{
// dgo has already finished the transaction; Commit again returns ErrFinished.
name: "does not commit again when CommitNow is set",
commitNow: true,
wantResp: true,
wantCommits: 0,
wantDiscards: 1,
},
{
name: "commit failure is returned",
txn: recordingTxn{commitErr: errCommitFailed},
wantErr: errCommitFailed,
wantCommits: 1,
wantDiscards: 1,
},
{
name: "mutate failure is returned without committing",
txn: recordingTxn{mutateErr: errMutationFailed},
wantErr: errMutationFailed,
wantCommits: 0,
wantDiscards: 1,
},
{
// A failing Discard must not turn a committed write into an error.
name: "discard failure does not fail a committed write",
txn: recordingTxn{discardErr: errDiscardFailed},
wantResp: true,
wantCommits: 1,
wantDiscards: 1,
},
}

for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
txn := tc.txn
client := setupWithTxn(t, &txn)

resp, err := client.Mutate(t.Context(), &api.Mutation{
SetJson: []byte(`{"name":"GoFr"}`),
CommitNow: tc.commitNow,
})

if tc.wantErr != nil {
require.ErrorIs(t, err, tc.wantErr)
require.Nil(t, resp)
} else {
require.NoError(t, err)
}

require.Equal(t, tc.wantResp, resp != nil, "response")
require.Equal(t, 1, txn.mutations, "Mutate calls")
require.Equal(t, tc.wantCommits, txn.commits, "Commit calls")
require.Equal(t, tc.wantDiscards, txn.discards, "Discard calls")
})
}
}

// Test_Mutate_DiscardSurvivesCallerCancellation pins that the deferred discard does not
// inherit the caller's cancellation.
//
// Discard runs after the commit has already returned. If it took the caller's context, a
// request canceled in that window would log "discard failed" for a write that was persisted
// -- an error line describing a success, which is the kind of log that sends someone looking
// for a data-loss bug that is not there.
func Test_Mutate_DiscardSurvivesCallerCancellation(t *testing.T) {
txn := recordingTxn{}
client := setupWithTxn(t, &txn)

ctx, cancel := context.WithCancel(t.Context())
cancel()

_, _ = client.Mutate(ctx, &api.Mutation{SetJson: []byte(`{"name":"GoFr"}`)})

require.Equal(t, 1, txn.discards, "the transaction is still discarded")
require.NoError(t, txn.discardCtxErr, "the discard must not inherit the caller's cancellation")
}
Loading