diff --git a/docs/datasources/dgraph/page.md b/docs/datasources/dgraph/page.md index c33b0927a3..20c1281597 100644 --- a/docs/datasources/dgraph/page.md +++ b/docs/datasources/dgraph/page.md @@ -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 diff --git a/pkg/gofr/container/datasources.go b/pkg/gofr/container/datasources.go index b929311798..b44b518162 100644 --- a/pkg/gofr/container/datasources.go +++ b/pkg/gofr/container/datasources.go @@ -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. diff --git a/pkg/gofr/datasource/dgraph/dgraph.go b/pkg/gofr/datasource/dgraph/dgraph.go index f884d613b0..d8b1f76ed3 100644 --- a/pkg/gofr/datasource/dgraph/dgraph.go +++ b/pkg/gofr/datasource/dgraph/dgraph.go @@ -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. +// +// 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() @@ -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 @@ -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 +// 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() diff --git a/pkg/gofr/datasource/dgraph/dgraph_test.go b/pkg/gofr/datasource/dgraph/dgraph_test.go index 150d26e6d7..664e3d171b 100644 --- a/pkg/gofr/datasource/dgraph/dgraph_test.go +++ b/pkg/gofr/datasource/dgraph/dgraph_test.go @@ -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()) @@ -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) diff --git a/pkg/gofr/datasource/dgraph/mutate_txn_test.go b/pkg/gofr/datasource/dgraph/mutate_txn_test.go new file mode 100644 index 0000000000..27e9e15109 --- /dev/null +++ b/pkg/gofr/datasource/dgraph/mutate_txn_test.go @@ -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") +}