Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
1fca5bd
fix: preserve prepared binary params across retries
ck89119 Jul 20, 2026
cf00b85
Merge branch 'main' into issue-25421-main
mergify[bot] Jul 20, 2026
3d22ec2
test: cover prepared binary retry and version boundary
ck89119 Jul 20, 2026
910cf1c
test: prove mixed-version binary dispatch boundary
ck89119 Jul 20, 2026
bc30225
Merge remote-tracking branch 'mo/main' into issue-25421-main
ck89119 Jul 20, 2026
0885830
Merge remote-tracking branch 'mo/main' into issue-25421-main
ck89119 Jul 20, 2026
ccba06c
Merge branch 'main' into issue-25421-main
ck89119 Jul 21, 2026
daca821
Merge remote-tracking branch 'mo/main' into issue-25421-main
ck89119 Jul 21, 2026
affe9b7
Merge remote-tracking branch 'mo/main' into issue-25421-main
ck89119 Jul 21, 2026
de8df52
fix: guard binary params in shard reads
ck89119 Jul 21, 2026
77b367e
fix: validate shard targets before remote reads
ck89119 Jul 21, 2026
d3d12d2
Merge remote-tracking branch 'mo/main' into issue-25421-main
ck89119 Jul 21, 2026
fa7e7b3
test: serialize shard read client instrumentation
ck89119 Jul 21, 2026
9e3f046
fix: preserve remote shard selection after validation
ck89119 Jul 21, 2026
e1e7d7e
fix: version binary shard read RPCs
ck89119 Jul 21, 2026
2be0d60
Merge remote-tracking branch 'mo/main' into issue-25421-main
ck89119 Jul 21, 2026
e0e549d
test: cover versioned shard read receivers
ck89119 Jul 21, 2026
abc51d4
Merge remote-tracking branch 'mo/main' into issue-25421-main
ck89119 Jul 21, 2026
6e185d3
Merge remote-tracking branch 'mo/main' into issue-25421-main
ck89119 Jul 22, 2026
daed241
fix: cancel rejected method RPC contexts
ck89119 Jul 22, 2026
fa8beb2
Merge branch 'main' into issue-25421-main
mergify[bot] Jul 22, 2026
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
1 change: 1 addition & 0 deletions pkg/common/morpc/method_based.go
Original file line number Diff line number Diff line change
Expand Up @@ -194,6 +194,7 @@ func (s *methodBasedServer[REQ, RESP]) onMessage(
resp := s.pool.AcquireResponse()
handlerCtx, ok := s.getHandleFunc(ctx, req, resp)
if !ok {
defer request.Cancel()
s.pool.ReleaseRequest(req)
return cs.Write(ctx, resp)
}
Expand Down
89 changes: 89 additions & 0 deletions pkg/common/morpc/method_based_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ package morpc
import (
"context"
"fmt"
"io"
"os"
"testing"
"time"
Expand All @@ -29,6 +30,94 @@ import (
"github.com/stretchr/testify/require"
)

type testMethodBasedClientSession struct {
write func(context.Context, Message) error
}

func (s *testMethodBasedClientSession) Close() error {
return nil
}

func (s *testMethodBasedClientSession) SessionCtx() context.Context {
return context.Background()
}

func (s *testMethodBasedClientSession) Write(ctx context.Context, message Message) error {
return s.write(ctx, message)
}

func (s *testMethodBasedClientSession) AsyncWrite(Message) error {
panic("not implemented")
}

func (s *testMethodBasedClientSession) CreateCache(context.Context, uint64) (MessageCache, error) {
panic("not implemented")
}

func (s *testMethodBasedClientSession) DeleteCache(uint64) {
panic("not implemented")
}

func (s *testMethodBasedClientSession) GetCache(uint64) (MessageCache, error) {
panic("not implemented")
}

func (s *testMethodBasedClientSession) RemoteAddress() string {
return ""
}

func TestMethodBasedServerCancelsRejectedRequest(t *testing.T) {
writeErr := io.ErrClosedPipe
for _, tc := range []struct {
name string
writeErr error
}{
{name: "write succeeds"},
{name: "write fails", writeErr: writeErr},
} {
t.Run(tc.name, func(t *testing.T) {
pool := NewMessagePool(
func() *testMethodBasedMessage { return &testMethodBasedMessage{} },
func() *testMethodBasedMessage { return &testMethodBasedMessage{} },
)
s := &methodBasedServer[*testMethodBasedMessage, *testMethodBasedMessage]{
logger: getLogger(""),
pool: pool,
handlers: make(map[uint32]handleFuncCtx[*testMethodBasedMessage, *testMethodBasedMessage]),
}

cancelCalls := 0
writeCalls := 0
cs := &testMethodBasedClientSession{
write: func(_ context.Context, message Message) error {
writeCalls++
require.Zero(t, cancelCalls)
require.True(t, moerr.IsMoErrCode(
message.(*testMethodBasedMessage).UnwrapError(),
moerr.ErrNotSupported,
))
return tc.writeErr
},
}
request := RPCMessage{
Message: &testMethodBasedMessage{method: 100},
Cancel: func() {
cancelCalls++
},
}

err := s.onMessage(t.Context(), request, 0, cs)
if tc.writeErr == nil {
require.NoError(t, err)
} else {
require.ErrorIs(t, err, tc.writeErr)
}
require.Equal(t, 1, writeCalls)
require.Equal(t, 1, cancelCalls)
})
}
}

func TestRPCSend(t *testing.T) {
runRPCTests(
t,
Expand Down
30 changes: 30 additions & 0 deletions pkg/container/vector/vector.go
Original file line number Diff line number Diff line change
Expand Up @@ -502,6 +502,36 @@ func NewVecWithData(
return vec
}

// NewVecWithDataCopy copies external backing data into allocations owned by mp.
func NewVecWithDataCopy(
typ types.Type,
length int,
data []byte,
area []byte,
mp *mpool.MPool,
) (*Vector, error) {
vec := NewVec(typ)
vec.length = length
var err error
if len(data) > 0 {
vec.data, err = mp.Alloc(len(data), false)
if err != nil {
vec.Free(mp)
return nil, err
}
copy(vec.data, data)
}
if len(area) > 0 {
vec.area, err = mp.Alloc(len(area), false)
if err != nil {
vec.Free(mp)
return nil, err
}
copy(vec.area, area)
}
return vec, nil
}

func NewConstNull(typ types.Type, length int, mp *mpool.MPool) *Vector {
vec := NewVecFromReuse()
vec.typ = typ
Expand Down
18 changes: 18 additions & 0 deletions pkg/container/vector/vector_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -441,6 +441,24 @@ func TestDup(t *testing.T) {
require.Equal(t, int64(0), mp.CurrNB())
}

func TestNewVecWithDataCopyOwnsBackingData(t *testing.T) {
mp := mpool.MustNewZero()
data := []byte("external-data")
area := []byte("external-area")
vec, err := NewVecWithDataCopy(types.T_text.ToType(), 1, data, area, mp)
require.NoError(t, err)
require.Equal(t, data, vec.GetData())
require.Equal(t, area, vec.GetArea())

data[0] = 'X'
area[0] = 'Y'
require.Equal(t, byte('e'), vec.GetData()[0])
require.Equal(t, byte('e'), vec.GetArea()[0])
require.NotPanics(t, func() { vec.Free(mp) })
require.Nil(t, vec.GetData())
require.Nil(t, vec.GetArea())
}

func TestShrink(t *testing.T) {
mp := mpool.MustNewZero()
{ // Array Float32
Expand Down
1 change: 1 addition & 0 deletions pkg/frontend/back_exec.go
Original file line number Diff line number Diff line change
Expand Up @@ -399,6 +399,7 @@ func doComQueryInBack(
proc.SetAffectedRows(backSes.lastAffectedRows)
proc.SetStmtProfile(&backSes.stmtProfile)
proc.SetResolveVariableFunc(backSes.txnCompileCtx.ResolveVariable)
proc.SetResolveVariableIsBinFunc(backSes.txnCompileCtx.ResolveVariableIsBin)
// Frontend back-exec — session-bound resolver. backSession is a
// frontend session without a client connection (NOT a system
// background task); all callers go through ses.GetBackgroundExec(...)
Expand Down
42 changes: 34 additions & 8 deletions pkg/frontend/compiler_context.go
Original file line number Diff line number Diff line change
Expand Up @@ -765,14 +765,8 @@ func (tcc *TxnCompilerContext) ResolveVariable(varName string, isSystemVar, isGl

ctx := tcc.execCtx.reqCtx

if ctx.Value(defines.InSp{}) != nil && ctx.Value(defines.InSp{}).(bool) {
tmpScope := ctx.Value(defines.VarScopeKey{}).(*[]map[string]interface{})
for i := len(*tmpScope) - 1; i >= 0; i-- {
curScope := (*tmpScope)[i]
if val, ok := curScope[strings.ToLower(varName)]; ok {
return val, nil
}
}
if val, ok := resolveStoredProcedureVariable(ctx, varName); ok {
return val, nil
}

if isSystemVar {
Expand All @@ -797,6 +791,38 @@ func (tcc *TxnCompilerContext) ResolveVariable(varName string, isSystemVar, isGl
return
}

func (tcc *TxnCompilerContext) ResolveVariableIsBin(varName string, isSystemVar, _ bool) (bool, error) {
if _, ok := resolveStoredProcedureVariable(tcc.execCtx.reqCtx, varName); ok {
return false, nil
}
if isSystemVar {
return false, nil
}
udVar, err := tcc.GetSession().GetUserDefinedVar(varName)
if err != nil {
return false, err
}
return udVar.IsBin, nil
}

func resolveStoredProcedureVariable(ctx context.Context, varName string) (interface{}, bool) {
inSp, _ := ctx.Value(defines.InSp{}).(bool)
if !inSp {
return nil, false
}
tmpScope, ok := ctx.Value(defines.VarScopeKey{}).(*[]map[string]interface{})
if !ok {
return nil, false
}
name := strings.ToLower(varName)
for i := len(*tmpScope) - 1; i >= 0; i-- {
if val, ok := (*tmpScope)[i][name]; ok {
return val, true
}
}
return nil, false
}

func (tcc *TxnCompilerContext) ResolveAccountIds(accountNames []string) (accountIds []uint32, err error) {
var sql string
var erArray []ExecResult
Expand Down
53 changes: 39 additions & 14 deletions pkg/frontend/computation_wrapper.go
Original file line number Diff line number Diff line change
Expand Up @@ -602,21 +602,11 @@ func initExecuteStmtParam(execCtx *ExecCtx, ses *Session, cwft *TxnComputationWr
if len(execPlan.Args) != numParams {
return nil, nil, nil, originSQL, moerr.NewInvalidInput(reqCtx, "Incorrect arguments to EXECUTE")
}
params := vector.NewVec(types.T_text.ToType())
paramVals := make([]any, numParams)
for i, arg := range execPlan.Args {
exprImpl := arg.Expr.(*plan.Expr_V)
param, err := cwft.proc.GetResolveVariableFunc()(exprImpl.V.Name, exprImpl.V.System, exprImpl.V.Global)
if err != nil {
return nil, nil, nil, originSQL, err
}
err = util.AppendAnyToStringVector(cwft.proc, param, params)
if err != nil {
return nil, nil, nil, originSQL, err
}
paramVals[i] = param
params, paramVals, paramIsBin, err := buildExecuteUserParams(cwft.proc, execPlan.Args)
if err != nil {
return nil, nil, nil, originSQL, err
}
cwft.proc.SetPrepareParams(params)
cwft.proc.SetOwnedPrepareParamsWithIsBin(params, paramIsBin)
cwft.paramVals = paramVals
} else {
if numParams > 0 {
Expand All @@ -626,6 +616,41 @@ func initExecuteStmtParam(execCtx *ExecCtx, ses *Session, cwft *TxnComputationWr
return prepareStmt.compile, preparePlan.Plan, prepareStmt.PrepareStmt, originSQL, nil
}

func buildExecuteUserParams(
proc *process.Process,
args []*plan.Expr,
) (params *vector.Vector, paramVals []any, paramIsBin []bool, err error) {
params = vector.NewVec(types.T_text.ToType())
defer func() {
if err != nil {
params.Free(proc.Mp())
}
}()
paramVals = make([]any, len(args))
paramIsBin = make([]bool, len(args))
for i, arg := range args {
exprImpl := arg.Expr.(*plan.Expr_V)
var param any
param, err = proc.GetResolveVariableFunc()(exprImpl.V.Name, exprImpl.V.System, exprImpl.V.Global)
if err != nil {
return
}
err = util.AppendAnyToStringVector(proc, param, params)
if err != nil {
return
}
resolveIsBin := proc.GetResolveVariableIsBinFunc()
if resolveIsBin != nil {
paramIsBin[i], err = resolveIsBin(exprImpl.V.Name, exprImpl.V.System, exprImpl.V.Global)
if err != nil {
return
}
}
paramVals[i] = plan2.ParamValue{Value: param, IsBin: paramIsBin[i]}
}
return
}

func shouldCachePrepareCompile(p *plan.Plan) bool {
if p == nil {
return true
Expand Down
Loading
Loading