MOT code refactoring

This commit is contained in:
Vinoth Veeraraghavan 2020-12-08 11:11:56 +08:00
parent 74de6dfc32
commit 282e48c468
23 changed files with 642 additions and 778 deletions

View File

@ -290,6 +290,8 @@ inline void Prefetch(const void* ptr)
#define MAX_KEY_COLUMNS (10U)
#define MAX_TUPLE_SIZE 16384 // in bytes
#define MAX_VARCHAR_LEN 1024
// Do not change this. Masstree assumes 15 for optimization purposes
#define BTREE_ORDER 15

View File

@ -71,6 +71,7 @@ bool RecoveryManager::RecoverDbStart()
MOT_LOG_INFO("Starting MOT recovery");
if (m_recoverFromCkptDone) {
SetLsn(m_lastReplayLsn);
return true;
}

View File

@ -26,7 +26,7 @@
#include <algorithm>
#include <unordered_map>
#include "../storage/table.h" // explicit path in order to solve collision with B db header file with the same name
#include "table.h"
#include "mot_engine.h"
#include "redo_log_writer.h"
#include "sentinel.h"

View File

@ -24,7 +24,7 @@
#include <map>
#include "../storage/table.h"
#include "table.h"
#include "row.h"
#include "txn.h"
#include "txn_access.h"

View File

@ -38,11 +38,11 @@ ENGINE_INC = $(top_builddir)/src/gausskernel/storage/mot/core/src
include $(top_builddir)/src/Makefile.global
OBJ_DIR = ../obj
OBJS = $(OBJ_DIR)/mot_fdw.o $(OBJ_DIR)/mot_internal.o $(OBJ_DIR)/mot_fdw_xlog.o $(OBJ_DIR)/mot_match_index.o
OBJS = $(OBJ_DIR)/mot_fdw.o $(OBJ_DIR)/mot_internal.o $(OBJ_DIR)/mot_fdw_xlog.o $(OBJ_DIR)/mot_match_index.o $(OBJ_DIR)/mot_fdw_error.o
DATA = mot_fdw.control mot_fdw--1.0.sql
DEPS := $(OBJ_DIR)/mot_fdw.d $(OBJ_DIR)/mot_internal.d $(OBJ_DIR)/mot_fdx_xlog.d $(OBJ_DIR)/mot_match_index.d
DEPS := $(OBJ_DIR)/mot_fdw.d $(OBJ_DIR)/mot_internal.d $(OBJ_DIR)/mot_fdx_xlog.d $(OBJ_DIR)/mot_match_index.d $(OBJ_DIR)/mot_fdw_error.d
# Shared library stuff
include $(top_srcdir)/src/gausskernel/common.mk

View File

@ -305,7 +305,7 @@ private:
memCfg.format("%" PRIu64 " MB", newLocalMemoryMb);
result = AddExtStringConfigItem("", "max_mot_local_memory", memCfg.c_str());
}
if (result) {
if (result && (motCfg.m_sessionLargeBufferStoreSizeMB != newSessionLargeStoreMemoryMb)) {
memCfg.format("%" PRIu64 " MB", newSessionLargeStoreMemoryMb);
result = AddExtStringConfigItem("", "session_large_buffer_store_size", memCfg.c_str());
}

View File

@ -41,7 +41,6 @@
#include "foreign/fdwapi.h"
#include "foreign/foreign.h"
#include "miscadmin.h"
#include "nodes/makefuncs.h"
#include "nodes/nodes.h"
#include "nodes/nodeFuncs.h"
#include "optimizer/cost.h"
@ -189,30 +188,6 @@ static int MOTGetFdwType()
return MOT_ORC;
}
static inline void BitmapDeSerialize(uint8_t* bitmap, int16_t len, ListCell** cell)
{
if (cell != nullptr && *cell != nullptr) {
int type = ((Const*)lfirst(*cell))->constvalue;
if (type == FDW_LIST_BITMAP) {
*cell = lnext(*cell);
for (int i = 0; i < len; i++) {
bitmap[i] = (uint8_t)((Const*)lfirst(*cell))->constvalue;
*cell = lnext(*cell);
}
}
}
}
static inline List* BitmapSerialize(List* result, uint8_t* bitmap, int16_t len)
{
// set list type to FDW_LIST_BITMAP
result = lappend(result, makeConst(INT4OID, -1, InvalidOid, 4, FDW_LIST_BITMAP, false, true));
for (int i = 0; i < len; i++)
result = lappend(result, makeConst(INT1OID, -1, InvalidOid, 1, Int8GetDatum(bitmap[i]), false, true));
return result;
}
void MOTRecover()
{
if (!MOTAdaptor::m_initialized) {
@ -1043,7 +1018,6 @@ static TupleTableSlot* MOTIterateForeignScan(ForeignScanState* node)
}
CleanQueryStatesOnError(festate->m_currTxn);
report_pg_error(rc,
festate->m_currTxn,
(void*)(festate->m_currTxn->m_errIx != nullptr ? festate->m_currTxn->m_errIx->GetName().c_str()
: "unknown"),
(void*)festate->m_currTxn->m_errMsgBuf);
@ -1091,7 +1065,6 @@ static TupleTableSlot* MOTIterateForeignScan(ForeignScanState* node)
CleanQueryStatesOnError(festate->m_currTxn);
report_pg_error(rc,
festate->m_currTxn,
(void*)(festate->m_currTxn->m_errIx != NULL ? festate->m_currTxn->m_errIx->GetName().c_str()
: "unknown"),
(void*)festate->m_currTxn->m_errMsgBuf);
@ -1418,7 +1391,7 @@ static TupleTableSlot* MOTExecForeignInsert(
fdwState->m_table = fdwState->m_currTxn->GetTableByExternalId(RelationGetRelid(resultRelInfo->ri_RelationDesc));
if (fdwState->m_table == nullptr) {
pfree(fdwState);
report_pg_error(MOT::RC_TABLE_NOT_FOUND, fdwState->m_currTxn);
report_pg_error(MOT::RC_TABLE_NOT_FOUND);
return nullptr;
}
fdwState->m_numAttrs = RelationGetNumberOfAttributes(resultRelInfo->ri_RelationDesc);
@ -1445,7 +1418,6 @@ static TupleTableSlot* MOTExecForeignInsert(
elog(DEBUG2, "Abort parent transaction from MOT insert, id %lu", fdwState->m_txnId);
CleanQueryStatesOnError(fdwState->m_currTxn);
report_pg_error(rc,
fdwState->m_currTxn,
(void*)(fdwState->m_currTxn->m_errIx != nullptr ? fdwState->m_currTxn->m_errIx->GetName().c_str()
: "unknown"),
(void*)fdwState->m_currTxn->m_errMsgBuf);
@ -1478,14 +1450,14 @@ static TupleTableSlot* MOTExecForeignUpdate(
planSlot->tts_nvalid,
(planSlot->tts_isnull[num] ? "NULL" : "NOT NULL"));
CleanQueryStatesOnError(fdwState->m_currTxn);
report_pg_error(MOT::RC_ERROR, fdwState->m_currTxn);
report_pg_error(MOT::RC_ERROR);
return nullptr;
}
if (currRow == nullptr) {
elog(ERROR, "MOTExecForeignUpdate failed to fetch row");
CleanQueryStatesOnError(fdwState->m_currTxn);
report_pg_error(((rc == MOT::RC_OK) ? MOT::RC_ERROR : rc), fdwState->m_currTxn);
report_pg_error(((rc == MOT::RC_OK) ? MOT::RC_ERROR : rc));
return nullptr;
}
@ -1504,7 +1476,6 @@ static TupleTableSlot* MOTExecForeignUpdate(
elog(DEBUG2, "Abort parent transaction from MOT update, id %lu", fdwState->m_txnId);
CleanQueryStatesOnError(fdwState->m_currTxn);
report_pg_error(rc,
fdwState->m_currTxn,
(void*)(fdwState->m_currTxn->m_errIx != nullptr ? fdwState->m_currTxn->m_errIx->GetName().c_str()
: "unknown"),
(void*)fdwState->m_currTxn->m_errMsgBuf);
@ -1532,7 +1503,7 @@ static TupleTableSlot* MOTExecForeignDelete(
planSlot->tts_nvalid,
(planSlot->tts_isnull[num] ? "NULL" : "NOT NULL"));
CleanQueryStatesOnError(fdwState->m_currTxn);
report_pg_error(MOT::RC_ERROR, fdwState->m_currTxn);
report_pg_error(MOT::RC_ERROR);
return nullptr;
}
@ -1543,7 +1514,7 @@ static TupleTableSlot* MOTExecForeignDelete(
}
elog(ERROR, "MOTExecForeignDelete failed to fetch row");
CleanQueryStatesOnError(fdwState->m_currTxn);
report_pg_error(((rc == MOT::RC_OK) ? MOT::RC_ERROR : rc), fdwState->m_currTxn);
report_pg_error(((rc == MOT::RC_OK) ? MOT::RC_ERROR : rc));
return nullptr;
}
@ -1561,7 +1532,6 @@ static TupleTableSlot* MOTExecForeignDelete(
elog(DEBUG2, "Abort parent transaction from MOT delete, id %lu", fdwState->m_txnId);
CleanQueryStatesOnError(fdwState->m_currTxn);
report_pg_error(rc,
fdwState->m_currTxn,
(void*)(fdwState->m_currTxn->m_errIx != nullptr ? fdwState->m_currTxn->m_errIx->GetName().c_str()
: "unknown"),
(void*)fdwState->m_currTxn->m_errMsgBuf);
@ -1700,7 +1670,7 @@ static void MOTXactCallback(XactEvent event, void* arg)
// commit this transaction. So we need to save the prepared transaction for further instructions from
// CN or gs_clean. If we failed to save it, it's a panic situation.
report_pg_error(
MOT::RC_PANIC, txn, (char*)"Failed to save prepared transaction data (FailedCommitPrepared)");
MOT::RC_PANIC, (char*)"Failed to save prepared transaction data (FailedCommitPrepared)");
}
return;
} else {
@ -1911,7 +1881,7 @@ static void MOTTruncateForeignTable(TruncateStmt* stmt, Relation rel)
errmodule(MOD_MOT),
errmsg("A checkpoint is in progress - cannot truncate table.")));
} else {
report_pg_error(rc, NULL, NULL, NULL, NULL, NULL, NULL);
report_pg_error(rc);
}
}
@ -1993,139 +1963,6 @@ static void InitMOTHandler()
}
}
MOTFdwStateSt* InitializeFdwState(void* fdwState, List** fdwExpr, uint64_t exTableID)
{
MOTFdwStateSt* state = (MOTFdwStateSt*)palloc0(sizeof(MOTFdwStateSt));
List* values = (List*)fdwState;
state->m_allocInScan = true;
state->m_foreignTableId = exTableID;
if (list_length(values) > 0) {
ListCell* cell = list_head(values);
int type = ((Const*)lfirst(cell))->constvalue;
if (type != FDW_LIST_STATE) {
return state;
}
cell = lnext(cell);
state->m_cmdOper = (CmdType)((Const*)lfirst(cell))->constvalue;
cell = lnext(cell);
state->m_order = (SORTDIR_ENUM)((Const*)lfirst(cell))->constvalue;
cell = lnext(cell);
state->m_hasForUpdate = (bool)((Const*)lfirst(cell))->constvalue;
cell = lnext(cell);
state->m_foreignTableId = ((Const*)lfirst(cell))->constvalue;
cell = lnext(cell);
state->m_numAttrs = ((Const*)lfirst(cell))->constvalue;
cell = lnext(cell);
state->m_ctidNum = ((Const*)lfirst(cell))->constvalue;
cell = lnext(cell);
state->m_numExpr = ((Const*)lfirst(cell))->constvalue;
cell = lnext(cell);
int len = BITMAP_GETLEN(state->m_numAttrs);
state->m_attrsUsed = (uint8_t*)palloc0(len);
state->m_attrsModified = (uint8_t*)palloc0(len);
BitmapDeSerialize(state->m_attrsUsed, len, &cell);
if (cell != NULL) {
state->m_bestIx = &state->m_bestIxBuf;
state->m_bestIx->Deserialize(cell, exTableID);
}
if (fdwExpr != NULL && *fdwExpr != NULL) {
ListCell* c = NULL;
int i = 0;
// divide fdw expr to param list and original expr
state->m_remoteCondsOrig = NULL;
foreach (c, *fdwExpr) {
if (i < state->m_numExpr) {
i++;
continue;
} else {
state->m_remoteCondsOrig = lappend(state->m_remoteCondsOrig, lfirst(c));
}
}
*fdwExpr = list_truncate(*fdwExpr, state->m_numExpr);
}
}
return state;
}
void* SerializeFdwState(MOTFdwStateSt* state)
{
List* result = NULL;
// set list type to FDW_LIST_STATE
result = lappend(result, makeConst(INT4OID, -1, InvalidOid, 4, FDW_LIST_STATE, false, true));
result = lappend(result, makeConst(INT4OID, -1, InvalidOid, 4, Int32GetDatum(state->m_cmdOper), false, true));
result = lappend(result, makeConst(INT1OID, -1, InvalidOid, 4, Int8GetDatum(state->m_order), false, true));
result = lappend(result, makeConst(BOOLOID, -1, InvalidOid, 1, BoolGetDatum(state->m_hasForUpdate), false, true));
result =
lappend(result, makeConst(INT4OID, -1, InvalidOid, 4, Int32GetDatum(state->m_foreignTableId), false, true));
result = lappend(result, makeConst(INT4OID, -1, InvalidOid, 4, Int32GetDatum(state->m_numAttrs), false, true));
result = lappend(result, makeConst(INT4OID, -1, InvalidOid, 4, Int32GetDatum(state->m_ctidNum), false, true));
result = lappend(result, makeConst(INT2OID, -1, InvalidOid, 2, Int16GetDatum(state->m_numExpr), false, true));
int len = BITMAP_GETLEN(state->m_numAttrs);
result = BitmapSerialize(result, state->m_attrsUsed, len);
if (state->m_bestIx != nullptr) {
state->m_bestIx->Serialize(&result);
}
ReleaseFdwState(state);
return result;
}
void CleanCursors(MOTFdwStateSt* state)
{
for (int i = 0; i < 2; i++) {
if (state->m_cursor[i]) {
state->m_cursor[i]->Invalidate();
state->m_cursor[i]->Destroy();
delete state->m_cursor[i];
state->m_cursor[i] = NULL;
}
}
}
void CleanQueryStatesOnError(MOT::TxnManager* txn)
{
if (txn != nullptr) {
for (auto& itr : txn->m_queryState) {
MOTFdwStateSt* state = (MOTFdwStateSt*)itr.second;
if (state != nullptr) {
CleanCursors(state);
}
}
}
}
void ReleaseFdwState(MOTFdwStateSt* state)
{
CleanCursors(state);
if (state->m_currTxn) {
state->m_currTxn->m_queryState.erase((uint64_t)state);
}
if (state->m_bestIx && state->m_bestIx != &state->m_bestIxBuf)
pfree(state->m_bestIx);
if (state->m_remoteCondsOrig != nullptr)
list_free(state->m_remoteCondsOrig);
if (state->m_attrsUsed != NULL)
pfree(state->m_attrsUsed);
if (state->m_attrsModified != NULL)
pfree(state->m_attrsModified);
state->m_table = NULL;
pfree(state);
}
void MOTCheckpointFetchLock()
{
MOT::MOTEngine* engine = MOT::MOTEngine::GetInstance();

View File

@ -0,0 +1,253 @@
/*
* Copyright (c) 2020 Huawei Technologies Co.,Ltd.
*
* openGauss is licensed under Mulan PSL v2.
* You can use this software according to the terms and conditions of the Mulan PSL v2.
* You may obtain a copy of Mulan PSL v2 at:
*
* http://license.coscl.org.cn/MulanPSL2
*
* THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY KIND,
* EITHER EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT,
* MERCHANTABILITY OR FIT FOR A PARTICULAR PURPOSE.
* See the Mulan PSL v2 for more details.
* -------------------------------------------------------------------------
*
* mot_fdw_error.cpp
* MOT Foreign Data Wrapper error reporting interfaces.
*
* IDENTIFICATION
* src/gausskernel/storage/mot/fdw_adapter/src/mot_fdw_error.cpp
*
* -------------------------------------------------------------------------
*/
#include "mot_fdw_error.h"
#include "catalog/namespace.h"
#include "nodes/parsenodes.h"
#include "utils/elog.h"
#include "column.h"
#include "table.h"
// Error code mapping array from MOT to PG
static const MotErrToPGErrSt MM_ERRCODE_TO_PG[] = {
// RC_OK
{ERRCODE_SUCCESSFUL_COMPLETION, "Success", nullptr},
// RC_ERROR
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_ABORT
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_UNSUPPORTED_COL_TYPE
{ERRCODE_INVALID_COLUMN_DEFINITION,
"Column definition of %s is not supported",
"Column type %s is not supported yet"},
// RC_UNSUPPORTED_COL_TYPE_ARR
{ERRCODE_INVALID_COLUMN_DEFINITION,
"Column definition of %s is not supported",
"Column type Array of %s is not supported yet"},
// RC_EXCEEDS_MAX_ROW_SIZE
{ERRCODE_FEATURE_NOT_SUPPORTED,
"Column definition of %s is not supported",
"Column size %d exceeds max tuple size %u"},
// RC_COL_NAME_EXCEEDS_MAX_SIZE
{ERRCODE_INVALID_COLUMN_DEFINITION,
"Column definition of %s is not supported",
"Column name %s exceeds max name size %u"},
// RC_COL_SIZE_INVALID
{ERRCODE_INVALID_COLUMN_DEFINITION,
"Column definition of %s is not supported",
"Column size %d exceeds max size %u"},
// RC_TABLE_EXCEEDS_MAX_DECLARED_COLS
{ERRCODE_FEATURE_NOT_SUPPORTED, "Can't create table", "Can't add column %s, number of declared columns is less"},
// RC_INDEX_EXCEEDS_MAX_SIZE
{ERRCODE_FDW_KEY_SIZE_EXCEEDS_MAX_ALLOWED,
"Can't create index",
"Total columns size is greater than maximum index size %u"},
// RC_TABLE_EXCEEDS_MAX_INDEXES,
{ERRCODE_FDW_TOO_MANY_INDEXES,
"Can't create index",
"Total number of indexes for table %s is greater than the maximum number if indexes allowed %u"},
// RC_TXN_EXCEEDS_MAX_DDLS,
{ERRCODE_FDW_TOO_MANY_DDL_CHANGES_IN_TRANSACTION_NOT_ALLOWED,
"Cannot execute statement",
"Maximum number of DDLs per transactions reached the maximum %u"},
// RC_UNIQUE_VIOLATION
{ERRCODE_UNIQUE_VIOLATION, "duplicate key value violates unique constraint \"%s\"", "Key %s already exists."},
// RC_TABLE_NOT_FOUND
{ERRCODE_UNDEFINED_TABLE, "Table \"%s\" doesn't exist", nullptr},
// RC_INDEX_NOT_FOUND
{ERRCODE_UNDEFINED_TABLE, "Index \"%s\" doesn't exist", nullptr},
// RC_LOCAL_ROW_FOUND
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_LOCAL_ROW_NOT_FOUND
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_LOCAL_ROW_DELETED
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_INSERT_ON_EXIST
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_INDEX_RETRY_INSERT
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_INDEX_DELETE
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_LOCAL_ROW_NOT_VISIBLE
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_MEMORY_ALLOCATION_ERROR
{ERRCODE_OUT_OF_LOGICAL_MEMORY, "Memory is temporarily unavailable", nullptr},
// RC_ILLEGAL_ROW_STATE
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_NULL_VOILATION
{ERRCODE_FDW_ERROR,
"Null constraint violated",
"NULL value cannot be inserted into non-null column %s at table %s"},
// RC_PANIC
{ERRCODE_FDW_ERROR, "Critical error", "Critical error: %s"},
// RC_NA
{ERRCODE_FDW_OPERATION_NOT_SUPPORTED, "A checkpoint is in progress - cannot truncate table.", nullptr},
// RC_MAX_VALUE
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr}
};
static_assert(sizeof(MM_ERRCODE_TO_PG) / sizeof(MotErrToPGErrSt) == MOT::RC_MAX_VALUE + 1,
"Not all MOT engine error codes (RC) is mapped to PG error codes");
void report_pg_error(MOT::RC rc, void* arg1, void* arg2, void* arg3, void* arg4, void* arg5)
{
const MotErrToPGErrSt* err = &MM_ERRCODE_TO_PG[rc];
switch (rc) {
case MOT::RC_OK:
break;
case MOT::RC_ERROR:
ereport(ERROR, (errmodule(MOD_MOT), errcode(err->m_pgErr), errmsg("%s", err->m_msg)));
break;
case MOT::RC_ABORT:
ereport(ERROR, (errmodule(MOD_MOT), errcode(err->m_pgErr), errmsg("%s", err->m_msg)));
break;
case MOT::RC_UNSUPPORTED_COL_TYPE: {
ColumnDef* col = (ColumnDef*)arg1;
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg(err->m_msg, col->colname),
errdetail(err->m_detail, NameListToString(col->typname->names))));
break;
}
case MOT::RC_UNSUPPORTED_COL_TYPE_ARR: {
ColumnDef* col = (ColumnDef*)arg1;
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg(err->m_msg, col->colname),
errdetail(err->m_detail, strVal(llast(col->typname->names)))));
break;
}
case MOT::RC_COL_NAME_EXCEEDS_MAX_SIZE: {
ColumnDef* col = (ColumnDef*)arg1;
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg(err->m_msg, col->colname),
errdetail(err->m_detail, col->colname, (uint32_t)MOT::Column::MAX_COLUMN_NAME_LEN)));
break;
}
case MOT::RC_COL_SIZE_INVALID: {
ColumnDef* col = (ColumnDef*)arg1;
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg(err->m_msg, col->colname),
errdetail(err->m_detail, (uint32_t)(uint64_t)arg2, (uint32_t)MAX_VARCHAR_LEN)));
break;
}
case MOT::RC_EXCEEDS_MAX_ROW_SIZE: {
ColumnDef* col = (ColumnDef*)arg1;
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg(err->m_msg, col->colname),
errdetail(err->m_detail, (uint32_t)(uint64_t)arg2, (uint32_t)MAX_TUPLE_SIZE)));
break;
}
case MOT::RC_TABLE_EXCEEDS_MAX_DECLARED_COLS: {
ColumnDef* col = (ColumnDef*)arg1;
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg("%s", err->m_msg),
errdetail(err->m_detail, col->colname)));
break;
}
case MOT::RC_INDEX_EXCEEDS_MAX_SIZE:
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg("%s", err->m_msg),
errdetail(err->m_detail, MAX_KEY_SIZE)));
break;
case MOT::RC_TABLE_EXCEEDS_MAX_INDEXES:
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg("%s", err->m_msg),
errdetail(err->m_detail, ((MOT::Table*)arg1)->GetTableName(), MAX_NUM_INDEXES)));
break;
case MOT::RC_TXN_EXCEEDS_MAX_DDLS:
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg("%s", err->m_msg),
errdetail(err->m_detail, MAX_DDL_ACCESS_SIZE)));
break;
case MOT::RC_UNIQUE_VIOLATION:
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg(err->m_msg, (char*)arg1),
errdetail(err->m_detail, (char*)arg2)));
break;
case MOT::RC_TABLE_NOT_FOUND:
case MOT::RC_INDEX_NOT_FOUND:
ereport(ERROR, (errmodule(MOD_MOT), errcode(err->m_pgErr), errmsg(err->m_msg, (char*)arg1)));
break;
// following errors are internal and should not get to an upper layer
case MOT::RC_LOCAL_ROW_FOUND:
case MOT::RC_LOCAL_ROW_NOT_FOUND:
case MOT::RC_LOCAL_ROW_DELETED:
case MOT::RC_INSERT_ON_EXIST:
case MOT::RC_INDEX_RETRY_INSERT:
case MOT::RC_INDEX_DELETE:
case MOT::RC_LOCAL_ROW_NOT_VISIBLE:
case MOT::RC_ILLEGAL_ROW_STATE:
ereport(ERROR, (errmodule(MOD_MOT), errcode(err->m_pgErr), errmsg("%s", err->m_msg)));
break;
case MOT::RC_MEMORY_ALLOCATION_ERROR:
ereport(ERROR, (errmodule(MOD_MOT), errcode(err->m_pgErr), errmsg("%s", err->m_msg)));
break;
case MOT::RC_NULL_VIOLATION: {
ColumnDef* col = (ColumnDef*)arg1;
MOT::Table* table = (MOT::Table*)arg2;
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg("%s", err->m_msg),
errdetail(err->m_detail, col->colname, table->GetLongTableName().c_str())));
break;
}
case MOT::RC_PANIC: {
char* msg = (char*)arg1;
ereport(PANIC,
(errmodule(MOD_MOT), errcode(err->m_pgErr), errmsg("%s", err->m_msg), errdetail(err->m_detail, msg)));
break;
}
case MOT::RC_NA:
ereport(ERROR, (errmodule(MOD_MOT), errcode(err->m_pgErr), errmsg("%s", err->m_msg)));
break;
case MOT::RC_MAX_VALUE:
default:
ereport(ERROR, (errmodule(MOD_MOT), errcode(err->m_pgErr), errmsg("%s", err->m_msg)));
break;
}
}

View File

@ -0,0 +1,39 @@
/*
* Copyright (c) 2020 Huawei Technologies Co.,Ltd.
*
* openGauss is licensed under Mulan PSL v2.
* You can use this software according to the terms and conditions of the Mulan PSL v2.
* You may obtain a copy of Mulan PSL v2 at:
*
* http://license.coscl.org.cn/MulanPSL2
*
* THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY KIND,
* EITHER EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT,
* MERCHANTABILITY OR FIT FOR A PARTICULAR PURPOSE.
* See the Mulan PSL v2 for more details.
* -------------------------------------------------------------------------
*
* mot_fdw_error.h
* MOT Foreign Data Wrapper error reporting interfaces.
*
* IDENTIFICATION
* src/gausskernel/storage/mot/fdw_adapter/src/mot_fdw_error.h
*
* -------------------------------------------------------------------------
*/
#ifndef MOT_FDW_ERROR_H
#define MOT_FDW_ERROR_H
#include "global.h"
typedef struct tagMotErrToPGErrSt {
int m_pgErr;
const char* m_msg;
const char* m_detail;
} MotErrToPGErrSt;
void report_pg_error(MOT::RC rc, void* arg1 = nullptr, void* arg2 = nullptr, void* arg3 = nullptr, void* arg4 = nullptr,
void* arg5 = nullptr);
#endif /* MOT_FDW_ERROR_H */

View File

@ -241,227 +241,6 @@ static void CancelSessionCleanup()
static GaussdbConfigLoader* gaussdbConfigLoader = nullptr;
// Error code mapping array from MOT to PG
static const MotErrToPGErrSt MM_ERRCODE_TO_PG[] = {
// RC_OK
{ERRCODE_SUCCESSFUL_COMPLETION, "Success", nullptr},
// RC_ERROR
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_ABORT
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_UNSUPPORTED_COL_TYPE
{ERRCODE_INVALID_COLUMN_DEFINITION,
"Column definition of %s is not supported",
"Column type %s is not supported yet"},
// RC_UNSUPPORTED_COL_TYPE_ARR
{ERRCODE_INVALID_COLUMN_DEFINITION,
"Column definition of %s is not supported",
"Column type Array of %s is not supported yet"},
// RC_EXCEEDS_MAX_ROW_SIZE
{ERRCODE_FEATURE_NOT_SUPPORTED,
"Column definition of %s is not supported",
"Column size %d exceeds max tuple size %u"},
// RC_COL_NAME_EXCEEDS_MAX_SIZE
{ERRCODE_INVALID_COLUMN_DEFINITION,
"Column definition of %s is not supported",
"Column name %s exceeds max name size %u"},
// RC_COL_SIZE_INVALID
{ERRCODE_INVALID_COLUMN_DEFINITION,
"Column definition of %s is not supported",
"Column size %d exceeds max size %u"},
// RC_TABLE_EXCEEDS_MAX_DECLARED_COLS
{ERRCODE_FEATURE_NOT_SUPPORTED, "Can't create table", "Can't add column %s, number of declared columns is less"},
// RC_INDEX_EXCEEDS_MAX_SIZE
{ERRCODE_FDW_KEY_SIZE_EXCEEDS_MAX_ALLOWED,
"Can't create index",
"Total columns size is greater than maximum index size %u"},
// RC_TABLE_EXCEEDS_MAX_INDEXES,
{ERRCODE_FDW_TOO_MANY_INDEXES,
"Can't create index",
"Total number of indexes for table %s is greater than the maximum number if indexes allowed %u"},
// RC_TXN_EXCEEDS_MAX_DDLS,
{ERRCODE_FDW_TOO_MANY_DDL_CHANGES_IN_TRANSACTION_NOT_ALLOWED,
"Cannot execute statement",
"Maximum number of DDLs per transactions reached the maximum %u"},
// RC_UNIQUE_VIOLATION
{ERRCODE_UNIQUE_VIOLATION, "duplicate key value violates unique constraint \"%s\"", "Key %s already exists."},
// RC_TABLE_NOT_FOUND
{ERRCODE_UNDEFINED_TABLE, "Table \"%s\" doesn't exist", nullptr},
// RC_INDEX_NOT_FOUND
{ERRCODE_UNDEFINED_TABLE, "Index \"%s\" doesn't exist", nullptr},
// RC_LOCAL_ROW_FOUND
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_LOCAL_ROW_NOT_FOUND
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_LOCAL_ROW_DELETED
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_INSERT_ON_EXIST
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_INDEX_RETRY_INSERT
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_INDEX_DELETE
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_LOCAL_ROW_NOT_VISIBLE
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_MEMORY_ALLOCATION_ERROR
{ERRCODE_OUT_OF_LOGICAL_MEMORY, "Memory is temporarily unavailable", nullptr},
// RC_ILLEGAL_ROW_STATE
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr},
// RC_NULL_VOILATION
{ERRCODE_FDW_ERROR,
"Null constraint violated",
"NULL value cannot be inserted into non-null column %s at table %s"},
// RC_PANIC
{ERRCODE_FDW_ERROR, "Critical error", "Critical error: %s"},
// RC_NA
{ERRCODE_FDW_OPERATION_NOT_SUPPORTED, "A checkpoint is in progress - cannot truncate table.", nullptr},
// RC_MAX_VALUE
{ERRCODE_FDW_ERROR, "Unknown error has occurred", nullptr}};
static_assert(sizeof(MM_ERRCODE_TO_PG) / sizeof(MotErrToPGErrSt) == MOT::RC_MAX_VALUE + 1,
"Not all MOT engine error codes (RC) is mapped to PG error codes");
void report_pg_error(MOT::RC rc, MOT::TxnManager* txn, void* arg1, void* arg2, void* arg3, void* arg4, void* arg5)
{
const MotErrToPGErrSt* err = &MM_ERRCODE_TO_PG[rc];
switch (rc) {
case MOT::RC_OK:
break;
case MOT::RC_ERROR:
ereport(ERROR, (errmodule(MOD_MOT), errcode(err->m_pgErr), errmsg("%s", err->m_msg)));
break;
case MOT::RC_ABORT:
ereport(ERROR, (errmodule(MOD_MOT), errcode(err->m_pgErr), errmsg("%s", err->m_msg)));
break;
case MOT::RC_UNSUPPORTED_COL_TYPE: {
ColumnDef* col = (ColumnDef*)arg1;
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg(err->m_msg, col->colname),
errdetail(err->m_detail, NameListToString(col->typname->names))));
break;
}
case MOT::RC_UNSUPPORTED_COL_TYPE_ARR: {
ColumnDef* col = (ColumnDef*)arg1;
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg(err->m_msg, col->colname),
errdetail(err->m_detail, strVal(llast(col->typname->names)))));
break;
}
case MOT::RC_COL_NAME_EXCEEDS_MAX_SIZE: {
ColumnDef* col = (ColumnDef*)arg1;
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg(err->m_msg, col->colname),
errdetail(err->m_detail, col->colname, (uint32_t)MOT::Column::MAX_COLUMN_NAME_LEN)));
break;
}
case MOT::RC_COL_SIZE_INVALID: {
ColumnDef* col = (ColumnDef*)arg1;
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg(err->m_msg, col->colname),
errdetail(err->m_detail, (uint32_t)(uint64_t)arg2, (uint32_t)MAX_VARCHAR_LEN)));
break;
}
case MOT::RC_EXCEEDS_MAX_ROW_SIZE: {
ColumnDef* col = (ColumnDef*)arg1;
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg(err->m_msg, col->colname),
errdetail(err->m_detail, (uint32_t)(uint64_t)arg2, (uint32_t)MAX_TUPLE_SIZE)));
break;
}
case MOT::RC_TABLE_EXCEEDS_MAX_DECLARED_COLS: {
ColumnDef* col = (ColumnDef*)arg1;
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg("%s", err->m_msg),
errdetail(err->m_detail, col->colname)));
break;
}
case MOT::RC_INDEX_EXCEEDS_MAX_SIZE:
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg("%s", err->m_msg),
errdetail(err->m_detail, MAX_KEY_SIZE)));
break;
case MOT::RC_TABLE_EXCEEDS_MAX_INDEXES:
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg("%s", err->m_msg),
errdetail(err->m_detail, ((MOT::Table*)arg1)->GetTableName(), MAX_NUM_INDEXES)));
break;
case MOT::RC_TXN_EXCEEDS_MAX_DDLS:
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg("%s", err->m_msg),
errdetail(err->m_detail, MAX_DDL_ACCESS_SIZE)));
break;
case MOT::RC_UNIQUE_VIOLATION:
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg(err->m_msg, (char*)arg1),
errdetail(err->m_detail, (char*)arg2)));
break;
case MOT::RC_TABLE_NOT_FOUND:
case MOT::RC_INDEX_NOT_FOUND:
ereport(ERROR, (errmodule(MOD_MOT), errcode(err->m_pgErr), errmsg(err->m_msg, (char*)arg1)));
break;
// following errors are internal and should not get to an upper layer
case MOT::RC_LOCAL_ROW_FOUND:
case MOT::RC_LOCAL_ROW_NOT_FOUND:
case MOT::RC_LOCAL_ROW_DELETED:
case MOT::RC_INSERT_ON_EXIST:
case MOT::RC_INDEX_RETRY_INSERT:
case MOT::RC_INDEX_DELETE:
case MOT::RC_LOCAL_ROW_NOT_VISIBLE:
case MOT::RC_ILLEGAL_ROW_STATE:
ereport(ERROR, (errmodule(MOD_MOT), errcode(err->m_pgErr), errmsg("%s", err->m_msg)));
break;
case MOT::RC_MEMORY_ALLOCATION_ERROR:
ereport(ERROR, (errmodule(MOD_MOT), errcode(err->m_pgErr), errmsg("%s", err->m_msg)));
break;
case MOT::RC_NULL_VIOLATION: {
ColumnDef* col = (ColumnDef*)arg1;
MOT::Table* table = (MOT::Table*)arg2;
ereport(ERROR,
(errmodule(MOD_MOT),
errcode(err->m_pgErr),
errmsg("%s", err->m_msg),
errdetail(err->m_detail, col->colname, table->GetLongTableName().c_str())));
break;
}
case MOT::RC_PANIC: {
char* msg = (char*)arg1;
ereport(PANIC,
(errmodule(MOD_MOT), errcode(err->m_pgErr), errmsg("%s", err->m_msg), errdetail(err->m_detail, msg)));
break;
}
case MOT::RC_NA:
ereport(ERROR, (errmodule(MOD_MOT), errcode(err->m_pgErr), errmsg("%s", err->m_msg)));
break;
case MOT::RC_MAX_VALUE:
default:
ereport(ERROR, (errmodule(MOD_MOT), errcode(err->m_pgErr), errmsg("%s", err->m_msg)));
break;
}
}
bool MOTAdaptor::m_initialized = false;
bool MOTAdaptor::m_callbacks_initialized = false;
@ -1328,7 +1107,7 @@ MOT::RC MOTAdaptor::CreateIndex(IndexStmt* index, ::TransactionId tid)
ix = MOT::IndexFactory::CreateIndex(index_order, indexing_method, flavor);
if (ix == nullptr) {
report_pg_error(MOT::RC_ABORT, txn);
report_pg_error(MOT::RC_ABORT);
return MOT::RC_ABORT;
}
ix->SetExtId(index->indexOid);
@ -1392,7 +1171,7 @@ MOT::RC MOTAdaptor::CreateIndex(IndexStmt* index, ::TransactionId tid)
if ((res = ix->IndexInit(keyLength, index->unique, index->idxname, nullptr)) != MOT::RC_OK) {
delete ix;
report_pg_error(res, txn);
report_pg_error(res);
return res;
}
@ -1406,7 +1185,7 @@ MOT::RC MOTAdaptor::CreateIndex(IndexStmt* index, ::TransactionId tid)
errmsg("Can not create index, max number of indexes %u reached", MAX_NUM_INDEXES)));
return MOT::RC_TABLE_EXCEEDS_MAX_INDEXES;
} else {
report_pg_error(txn->m_err, txn, index->idxname, txn->m_errMsgBuf);
report_pg_error(txn->m_err, index->idxname, txn->m_errMsgBuf);
return MOT::RC_UNIQUE_VIOLATION;
}
}
@ -1464,7 +1243,7 @@ MOT::RC MOTAdaptor::CreateTable(CreateForeignTableStmt* table, ::TransactionId t
table->base.relation->relname, tname.c_str(), columnCount, table->base.relation->foreignOid)) {
delete currentTable;
currentTable = nullptr;
report_pg_error(MOT::RC_MEMORY_ALLOCATION_ERROR, txn);
report_pg_error(MOT::RC_MEMORY_ALLOCATION_ERROR);
break;
}
@ -1474,7 +1253,7 @@ MOT::RC MOTAdaptor::CreateTable(CreateForeignTableStmt* table, ::TransactionId t
if (res != MOT::RC_OK) {
delete currentTable;
currentTable = nullptr;
report_pg_error(MOT::RC_MEMORY_ALLOCATION_ERROR, txn);
report_pg_error(MOT::RC_MEMORY_ALLOCATION_ERROR);
break;
}
@ -1501,7 +1280,7 @@ MOT::RC MOTAdaptor::CreateTable(CreateForeignTableStmt* table, ::TransactionId t
if (res != MOT::RC_OK) {
delete currentTable;
currentTable = nullptr;
report_pg_error(res, txn, colDef, (void*)(int64)typeLen);
report_pg_error(res, colDef, (void*)(int64)typeLen);
break;
}
hasBlob |= isBlob;
@ -1555,7 +1334,7 @@ MOT::RC MOTAdaptor::CreateTable(CreateForeignTableStmt* table, ::TransactionId t
if (res != MOT::RC_OK) {
delete currentTable;
currentTable = nullptr;
report_pg_error(res, txn, colDef, (void*)(int64)typeLen);
report_pg_error(res, colDef, (void*)(int64)typeLen);
break;
}
}
@ -1583,7 +1362,7 @@ MOT::RC MOTAdaptor::CreateTable(CreateForeignTableStmt* table, ::TransactionId t
if (!currentTable->InitRowPool()) {
delete currentTable;
currentTable = nullptr;
report_pg_error(MOT::RC_MEMORY_ALLOCATION_ERROR, txn);
report_pg_error(MOT::RC_MEMORY_ALLOCATION_ERROR);
break;
}
@ -1598,7 +1377,7 @@ MOT::RC MOTAdaptor::CreateTable(CreateForeignTableStmt* table, ::TransactionId t
if (res != MOT::RC_OK) {
delete currentTable;
currentTable = nullptr;
report_pg_error(res, txn);
report_pg_error(res);
break;
}
@ -1613,7 +1392,7 @@ MOT::RC MOTAdaptor::CreateTable(CreateForeignTableStmt* table, ::TransactionId t
if (rc != MOT::RC_OK) {
delete currentTable;
currentTable = nullptr;
report_pg_error(rc, txn);
report_pg_error(rc);
break;
}
primaryIdx->SetExtId(table->base.relation->foreignOid + 1);
@ -2331,3 +2110,112 @@ void MOTAdaptor::DatumToMOTKey(
break;
}
}
MOTFdwStateSt* InitializeFdwState(void* fdwState, List** fdwExpr, uint64_t exTableID)
{
MOTFdwStateSt* state = (MOTFdwStateSt*)palloc0(sizeof(MOTFdwStateSt));
List* values = (List*)fdwState;
state->m_allocInScan = true;
state->m_foreignTableId = exTableID;
if (list_length(values) > 0) {
ListCell* cell = list_head(values);
int type = ((Const*)lfirst(cell))->constvalue;
if (type != FDW_LIST_STATE) {
return state;
}
cell = lnext(cell);
state->m_cmdOper = (CmdType)((Const*)lfirst(cell))->constvalue;
cell = lnext(cell);
state->m_order = (SORTDIR_ENUM)((Const*)lfirst(cell))->constvalue;
cell = lnext(cell);
state->m_hasForUpdate = (bool)((Const*)lfirst(cell))->constvalue;
cell = lnext(cell);
state->m_foreignTableId = ((Const*)lfirst(cell))->constvalue;
cell = lnext(cell);
state->m_numAttrs = ((Const*)lfirst(cell))->constvalue;
cell = lnext(cell);
state->m_ctidNum = ((Const*)lfirst(cell))->constvalue;
cell = lnext(cell);
state->m_numExpr = ((Const*)lfirst(cell))->constvalue;
cell = lnext(cell);
int len = BITMAP_GETLEN(state->m_numAttrs);
state->m_attrsUsed = (uint8_t*)palloc0(len);
state->m_attrsModified = (uint8_t*)palloc0(len);
BitmapDeSerialize(state->m_attrsUsed, len, &cell);
if (cell != NULL) {
state->m_bestIx = &state->m_bestIxBuf;
state->m_bestIx->Deserialize(cell, exTableID);
}
if (fdwExpr != NULL && *fdwExpr != NULL) {
ListCell* c = NULL;
int i = 0;
// divide fdw expr to param list and original expr
state->m_remoteCondsOrig = NULL;
foreach (c, *fdwExpr) {
if (i < state->m_numExpr) {
i++;
continue;
} else {
state->m_remoteCondsOrig = lappend(state->m_remoteCondsOrig, lfirst(c));
}
}
*fdwExpr = list_truncate(*fdwExpr, state->m_numExpr);
}
}
return state;
}
void* SerializeFdwState(MOTFdwStateSt* state)
{
List* result = NULL;
// set list type to FDW_LIST_STATE
result = lappend(result, makeConst(INT4OID, -1, InvalidOid, 4, FDW_LIST_STATE, false, true));
result = lappend(result, makeConst(INT4OID, -1, InvalidOid, 4, Int32GetDatum(state->m_cmdOper), false, true));
result = lappend(result, makeConst(INT1OID, -1, InvalidOid, 4, Int8GetDatum(state->m_order), false, true));
result = lappend(result, makeConst(BOOLOID, -1, InvalidOid, 1, BoolGetDatum(state->m_hasForUpdate), false, true));
result =
lappend(result, makeConst(INT4OID, -1, InvalidOid, 4, Int32GetDatum(state->m_foreignTableId), false, true));
result = lappend(result, makeConst(INT4OID, -1, InvalidOid, 4, Int32GetDatum(state->m_numAttrs), false, true));
result = lappend(result, makeConst(INT4OID, -1, InvalidOid, 4, Int32GetDatum(state->m_ctidNum), false, true));
result = lappend(result, makeConst(INT2OID, -1, InvalidOid, 2, Int16GetDatum(state->m_numExpr), false, true));
int len = BITMAP_GETLEN(state->m_numAttrs);
result = BitmapSerialize(result, state->m_attrsUsed, len);
if (state->m_bestIx != nullptr) {
state->m_bestIx->Serialize(&result);
}
ReleaseFdwState(state);
return result;
}
void ReleaseFdwState(MOTFdwStateSt* state)
{
CleanCursors(state);
if (state->m_currTxn) {
state->m_currTxn->m_queryState.erase((uint64_t)state);
}
if (state->m_bestIx && state->m_bestIx != &state->m_bestIxBuf)
pfree(state->m_bestIx);
if (state->m_remoteCondsOrig != nullptr)
list_free(state->m_remoteCondsOrig);
if (state->m_attrsUsed != NULL)
pfree(state->m_attrsUsed);
if (state->m_attrsModified != NULL)
pfree(state->m_attrsModified);
state->m_table = NULL;
pfree(state);
}

View File

@ -22,18 +22,20 @@
* -------------------------------------------------------------------------
*/
#ifndef MOT_INT_H
#define MOT_INT_H
#ifndef MOT_INTERNAL_H
#define MOT_INTERNAL_H
#include <map>
#include <string>
#include "catalog_column_types.h"
#include "foreign/fdwapi.h"
#include "nodes/nodes.h"
#include "nodes/makefuncs.h"
#include "utils/numeric.h"
#include "utils/numeric_gs.h"
#include "pgstat.h"
#include "global.h"
#include "mot_fdw_error.h"
#include "mot_fdw_xlog.h"
#include "system/mot_engine.h"
#include "bitmapset.h"
@ -92,17 +94,6 @@ class MOTEngine;
typedef struct MOTFdwState_St MOTFdwStateSt;
#endif
typedef struct MotErrToPGErr_St {
int m_pgErr;
const char* m_msg;
const char* m_detail;
} MotErrToPGErrSt;
void report_pg_error(MOT::RC rc, MOT::TxnManager* txn, void* arg1 = nullptr, void* arg2 = nullptr, void* arg3 = nullptr,
void* arg4 = nullptr, void* arg5 = nullptr);
#define MAX_VARCHAR_LEN 1024
typedef enum : uint8_t { SORTDIR_NONE = 0, SORTDIR_ASC = 1, SORTDIR_DESC = 2 } SORTDIR_ENUM;
typedef enum : uint8_t { FDW_LIST_STATE = 1, FDW_LIST_BITMAP = 2 } FDW_LIST_TYPE;
@ -332,7 +323,7 @@ inline MOT::TxnManager* GetSafeTxn(const char* callerSrc, ::TransactionId txn_id
u_sess->mot_cxt.txn_manager->SetTransactionId(txn_id);
}
} else {
report_pg_error(MOT_GET_ROOT_ERROR_RC(), nullptr, nullptr, nullptr, nullptr, nullptr, nullptr);
report_pg_error(MOT_GET_ROOT_ERROR_RC());
}
}
return u_sess->mot_cxt.txn_manager;
@ -340,4 +331,56 @@ inline MOT::TxnManager* GetSafeTxn(const char* callerSrc, ::TransactionId txn_id
extern void EnsureSafeThreadAccess();
#endif // MOT_INT_H
inline List* BitmapSerialize(List* result, uint8_t* bitmap, int16_t len)
{
// set list type to FDW_LIST_BITMAP
result = lappend(result, makeConst(INT4OID, -1, InvalidOid, 4, FDW_LIST_BITMAP, false, true));
for (int i = 0; i < len; i++)
result = lappend(result, makeConst(INT1OID, -1, InvalidOid, 1, Int8GetDatum(bitmap[i]), false, true));
return result;
}
inline void BitmapDeSerialize(uint8_t* bitmap, int16_t len, ListCell** cell)
{
if (cell != nullptr && *cell != nullptr) {
int type = ((Const*)lfirst(*cell))->constvalue;
if (type == FDW_LIST_BITMAP) {
*cell = lnext(*cell);
for (int i = 0; i < len; i++) {
bitmap[i] = (uint8_t)((Const*)lfirst(*cell))->constvalue;
*cell = lnext(*cell);
}
}
}
}
inline void CleanCursors(MOTFdwStateSt* state)
{
for (int i = 0; i < 2; i++) {
if (state->m_cursor[i]) {
state->m_cursor[i]->Invalidate();
state->m_cursor[i]->Destroy();
delete state->m_cursor[i];
state->m_cursor[i] = NULL;
}
}
}
inline void CleanQueryStatesOnError(MOT::TxnManager* txn)
{
if (txn != nullptr) {
for (auto& itr : txn->m_queryState) {
MOTFdwStateSt* state = (MOTFdwStateSt*)itr.second;
if (state != nullptr) {
CleanCursors(state);
}
}
}
}
MOTFdwStateSt* InitializeFdwState(void* fdwState, List** fdwExpr, uint64_t exTableID);
void* SerializeFdwState(MOTFdwStateSt* state);
void ReleaseFdwState(MOTFdwStateSt* state);
#endif // MOT_INTERNAL_H

View File

@ -55,11 +55,9 @@
#include "mot_engine.h"
#include "utilities.h"
#include "mot_internal.h"
#include "catalog_column_types.h"
#include "mot_error.h"
#include "utilities.h"
#include "mot_internal.h"
#include "cycles.h"
#include <assert.h>
@ -195,7 +193,7 @@ static void ProcessJitResult(MOT::RC result, JitContext* jitContext, int newScan
} else {
JitStatisticsProvider::GetInstance().AddFailExecQuery();
}
report_pg_error(result, currTxn, (void*)arg1, (void*)arg2);
report_pg_error(result, (void*)arg1, (void*)arg2);
}
}
@ -369,7 +367,7 @@ extern int JitExecQuery(
"Execute JIT",
"Cannot execute jitted function: function is null. Aborting transaction.");
JitStatisticsProvider::GetInstance().AddFailExecQuery();
report_pg_error(MOT::RC_ERROR, NULL); // execution control ends, calls ereport(error,...)
report_pg_error(MOT::RC_ERROR); // execution control ends, calls ereport(error,...)
}
// when running under thread-pool, it is possible to be executed from a thread that hasn't yet executed any MOT code
@ -384,7 +382,7 @@ extern int JitExecQuery(
"Execute JIT",
"Cannot execute jitted code: Current transaction is undefined. Aborting transaction.");
JitStatisticsProvider::GetInstance().AddFailExecQuery();
report_pg_error(MOT::RC_MEMORY_ALLOCATION_ERROR, NULL); // execution control ends, calls ereport(error,...)
report_pg_error(MOT::RC_MEMORY_ALLOCATION_ERROR); // execution control ends, calls ereport(error,...)
}
// during the very first invocation of the query we need to setup the reusable search keys
@ -395,7 +393,7 @@ extern int JitExecQuery(
MOT_REPORT_ERROR(
MOT_ERROR_OOM, "Execute JIT", "Failed to prepare for executing jitted code, aborting transaction");
JitStatisticsProvider::GetInstance().AddFailExecQuery();
report_pg_error(MOT::RC_MEMORY_ALLOCATION_ERROR, NULL); // execution control ends, calls ereport(error,...)
report_pg_error(MOT::RC_MEMORY_ALLOCATION_ERROR); // execution control ends, calls ereport(error,...)
}
}

View File

@ -22,8 +22,13 @@
* -------------------------------------------------------------------------
*/
#include "jit_llvm_funcs.h"
/*
* ATTENTION: Be sure to include jit_llvm_query.h before anything else because of gscodegen.h
* (jit_llvm_blocks.h includes jit_llvm_query.h before anything else).
* See jit_llvm_query.h for more details.
*/
#include "jit_llvm_blocks.h"
#include "jit_llvm_funcs.h"
#include "jit_util.h"
#include "mot_error.h"
#include "utilities.h"
@ -35,6 +40,15 @@ using namespace dorado;
namespace JitExec {
DECLARE_LOGGER(JitLlvmBlocks, JitExec)
static bool ProcessJoinOpExpr(
JitLlvmCodeGenContext* ctx, const OpExpr* op_expr, int* column_count, int* column_array, int* max_arg);
static bool ProcessJoinBoolExpr(
JitLlvmCodeGenContext* ctx, const BoolExpr* boolexpr, int* column_count, int* column_array, int* max_arg);
static llvm::Value* ProcessFilterExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitFilter* filter, int* max_arg);
static llvm::Value* ProcessExpr(
JitLlvmCodeGenContext* ctx, Expr* expr, int& result_type, int arg_pos, int depth, int* max_arg);
static llvm::Value* ProcessExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitExpr* expr, int* max_arg);
/*--------------------------- Helpers to generate compound LLVM code ---------------------------*/
/** @brief Builds a code segment for checking if soft memory limit has been reached. */
void buildIsSoftMemoryLimitReached(JitLlvmCodeGenContext* ctx)
@ -50,7 +64,7 @@ void buildIsSoftMemoryLimitReached(JitLlvmCodeGenContext* ctx)
}
/** @brief Builds a code segment for writing datum value to a column. */
void buildWriteDatumColumn(JitLlvmCodeGenContext* ctx, llvm::Value* row, int colid, llvm::Value* datum_value)
static void buildWriteDatumColumn(JitLlvmCodeGenContext* ctx, llvm::Value* row, int colid, llvm::Value* datum_value)
{
llvm::Value* set_null_bit_res = AddSetExprResultNullBit(ctx, row, colid);
IssueDebugLog("Set null bit");
@ -90,7 +104,7 @@ void buildWriteRow(JitLlvmCodeGenContext* ctx, llvm::Value* row, bool isPKey, Ji
}
/** @brief Process a join expression (WHERE clause) and generate code to build a search key. */
bool ProcessJoinExpr(JitLlvmCodeGenContext* ctx, Expr* expr, int* column_count, int* column_array, int* max_arg)
static bool ProcessJoinExpr(JitLlvmCodeGenContext* ctx, Expr* expr, int* column_count, int* column_array, int* max_arg)
{
bool result = false;
if (expr->type == T_OpExpr) {
@ -104,7 +118,7 @@ bool ProcessJoinExpr(JitLlvmCodeGenContext* ctx, Expr* expr, int* column_count,
}
/** @brief Process an operator expression (process only "COLUMN equals EXPR" operators). */
bool ProcessJoinOpExpr(
static bool ProcessJoinOpExpr(
JitLlvmCodeGenContext* ctx, const OpExpr* op_expr, int* column_count, int* column_array, int* max_arg)
{
bool result = false;
@ -176,7 +190,7 @@ bool ProcessJoinOpExpr(
/** @brief Process a boolean operator (process only AND operators, since we handle only point queries, or full-prefix
* range update). */
bool ProcessJoinBoolExpr(
static bool ProcessJoinBoolExpr(
JitLlvmCodeGenContext* ctx, const BoolExpr* boolexpr, int* column_count, int* column_array, int* max_arg)
{
bool result = false;
@ -294,7 +308,7 @@ llvm::Value* buildSearchRow(JitLlvmCodeGenContext* ctx, MOT::AccessType access_t
return row;
}
llvm::Value* buildFilter(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitFilter* filter, int* max_arg)
static llvm::Value* buildFilter(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitFilter* filter, int* max_arg)
{
llvm::Value* result = nullptr;
llvm::Value* lhs_expr = ProcessExpr(ctx, row, filter->_lhs_operand, max_arg);
@ -370,8 +384,8 @@ void buildDeleteRow(JitLlvmCodeGenContext* ctx)
}
/** @brief Adds code to search for an iterator. */
llvm::Value* buildSearchIterator(JitLlvmCodeGenContext* ctx, JitIndexScanDirection index_scan_direction,
JitRangeBoundMode range_bound_mode, JitRangeScanType range_scan_type, int subQueryIndex /* = -1 */)
static llvm::Value* buildSearchIterator(JitLlvmCodeGenContext* ctx, JitIndexScanDirection index_scan_direction,
JitRangeBoundMode range_bound_mode, JitRangeScanType range_scan_type, int subQueryIndex = -1)
{
// search the row
IssueDebugLog("Searching range start");
@ -388,8 +402,8 @@ llvm::Value* buildSearchIterator(JitLlvmCodeGenContext* ctx, JitIndexScanDirecti
}
/** @brief Adds code to search for an iterator. */
llvm::Value* buildBeginIterator(
JitLlvmCodeGenContext* ctx, JitRangeScanType rangeScanType, int subQueryIndex /* = -1 */)
static llvm::Value* buildBeginIterator(
JitLlvmCodeGenContext* ctx, JitRangeScanType rangeScanType, int subQueryIndex = -1)
{
// search the row
IssueDebugLog("Getting begin iterator for full-scan");
@ -426,7 +440,7 @@ llvm::Value* buildGetRowFromIterator(JitLlvmCodeGenContext* ctx, llvm::BasicBloc
}
/** @brief Process constant expression. */
llvm::Value* ProcessConstExpr(
static llvm::Value* ProcessConstExpr(
JitLlvmCodeGenContext* ctx, const Const* const_value, int& result_type, int arg_pos, int depth, int* max_arg)
{
llvm::Value* result = nullptr;
@ -451,7 +465,7 @@ llvm::Value* ProcessConstExpr(
}
/** @brief Process Param expression. */
llvm::Value* ProcessParamExpr(
static llvm::Value* ProcessParamExpr(
JitLlvmCodeGenContext* ctx, const Param* param, int& result_type, int arg_pos, int depth, int* max_arg)
{
llvm::Value* result = nullptr;
@ -475,7 +489,7 @@ llvm::Value* ProcessParamExpr(
}
/** @brief Process Relabel expression as Param expression. */
llvm::Value* ProcessRelabelExpr(
static llvm::Value* ProcessRelabelExpr(
JitLlvmCodeGenContext* ctx, RelabelType* relabel_type, int& result_type, int arg_pos, int depth, int* max_arg)
{
llvm::Value* result = nullptr;
@ -491,7 +505,7 @@ llvm::Value* ProcessRelabelExpr(
}
/** @brief Proess Var expression. */
llvm::Value* ProcessVarExpr(
static llvm::Value* ProcessVarExpr(
JitLlvmCodeGenContext* ctx, const Var* var, int& result_type, int arg_pos, int depth, int* max_arg)
{
llvm::Value* result = nullptr;
@ -516,7 +530,7 @@ llvm::Value* ProcessVarExpr(
}
/** @brief Adds call to PG unary operator. */
llvm::Value* AddExecUnaryOperator(
static llvm::Value* AddExecUnaryOperator(
JitLlvmCodeGenContext* ctx, llvm::Value* param, llvm::Constant* unary_operator, int arg_pos)
{
llvm::Constant* arg_pos_value = llvm::ConstantInt::get(ctx->INT32_T, arg_pos, true);
@ -524,7 +538,7 @@ llvm::Value* AddExecUnaryOperator(
}
/** @brief Adds call to PG binary operator. */
llvm::Value* AddExecBinaryOperator(JitLlvmCodeGenContext* ctx, llvm::Value* lhs_param, llvm::Value* rhs_param,
static llvm::Value* AddExecBinaryOperator(JitLlvmCodeGenContext* ctx, llvm::Value* lhs_param, llvm::Value* rhs_param,
llvm::Constant* binary_operator, int arg_pos)
{
llvm::Constant* arg_pos_value = llvm::ConstantInt::get(ctx->INT32_T, arg_pos, true);
@ -532,7 +546,7 @@ llvm::Value* AddExecBinaryOperator(JitLlvmCodeGenContext* ctx, llvm::Value* lhs_
}
/** @brief Adds call to PG ternary operator. */
llvm::Value* AddExecTernaryOperator(JitLlvmCodeGenContext* ctx, llvm::Value* param1, llvm::Value* param2,
static llvm::Value* AddExecTernaryOperator(JitLlvmCodeGenContext* ctx, llvm::Value* param1, llvm::Value* param2,
llvm::Value* param3, llvm::Constant* ternary_operator, int arg_pos)
{
llvm::Constant* arg_pos_value = llvm::ConstantInt::get(ctx->INT32_T, arg_pos, true);
@ -562,7 +576,7 @@ llvm::Value* AddExecTernaryOperator(JitLlvmCodeGenContext* ctx, llvm::Value* par
#define APPLY_TERNARY_CAST_OPERATOR(funcid, name) APPLY_TERNARY_OPERATOR(funcid, name)
/** @brief Process operator expression. */
llvm::Value* ProcessOpExpr(
static llvm::Value* ProcessOpExpr(
JitLlvmCodeGenContext* ctx, const OpExpr* op_expr, int& result_type, int arg_pos, int depth, int* max_arg)
{
llvm::Value* result = nullptr;
@ -610,7 +624,7 @@ llvm::Value* ProcessOpExpr(
}
/** @brief Process function expression. */
llvm::Value* ProcessFuncExpr(
static llvm::Value* ProcessFuncExpr(
JitLlvmCodeGenContext* ctx, const FuncExpr* func_expr, int& result_type, int arg_pos, int depth, int* max_arg)
{
llvm::Value* result = nullptr;
@ -683,7 +697,7 @@ llvm::Value* ProcessFuncExpr(
MOT_LOG_TRACE("Unexpected call in filter expression to ternary cast builtin: " #name); \
break;
llvm::Value* ProcessFilterExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitFilter* filter, int* max_arg)
static llvm::Value* ProcessFilterExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitFilter* filter, int* max_arg)
{
llvm::Value* result = nullptr;
@ -721,7 +735,8 @@ llvm::Value* ProcessFilterExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, Jit
#undef APPLY_TERNARY_CAST_OPERATOR
/** @brief Process an expression. Generates code to evaluate the expression. */
llvm::Value* ProcessExpr(JitLlvmCodeGenContext* ctx, Expr* expr, int& result_type, int arg_pos, int depth, int* max_arg)
static llvm::Value* ProcessExpr(
JitLlvmCodeGenContext* ctx, Expr* expr, int& result_type, int arg_pos, int depth, int* max_arg)
{
llvm::Value* result = nullptr;
MOT_LOG_DEBUG("%*s --> Processing expression %d", depth, "", (int)expr->type);
@ -752,7 +767,7 @@ llvm::Value* ProcessExpr(JitLlvmCodeGenContext* ctx, Expr* expr, int& result_typ
return result;
}
llvm::Value* ProcessConstExpr(JitLlvmCodeGenContext* ctx, const JitConstExpr* expr, int* max_arg)
static llvm::Value* ProcessConstExpr(JitLlvmCodeGenContext* ctx, const JitConstExpr* expr, int* max_arg)
{
AddSetExprArgIsNull(ctx, expr->_arg_pos, expr->_is_null); // mark expression null status
llvm::Value* result = llvm::ConstantInt::get(ctx->INT64_T, expr->_value, true);
@ -762,7 +777,7 @@ llvm::Value* ProcessConstExpr(JitLlvmCodeGenContext* ctx, const JitConstExpr* ex
return result;
}
llvm::Value* ProcessParamExpr(JitLlvmCodeGenContext* ctx, const JitParamExpr* expr, int* max_arg)
static llvm::Value* ProcessParamExpr(JitLlvmCodeGenContext* ctx, const JitParamExpr* expr, int* max_arg)
{
llvm::Value* result = AddGetDatumParam(ctx, expr->_param_id, expr->_arg_pos);
if (max_arg && (expr->_arg_pos > *max_arg)) {
@ -771,7 +786,7 @@ llvm::Value* ProcessParamExpr(JitLlvmCodeGenContext* ctx, const JitParamExpr* ex
return result;
}
llvm::Value* ProcessVarExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, const JitVarExpr* expr, int* max_arg)
static llvm::Value* ProcessVarExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, const JitVarExpr* expr, int* max_arg)
{
llvm::Value* result = nullptr;
if (row == nullptr) {
@ -809,7 +824,7 @@ llvm::Value* ProcessVarExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, const
#define APPLY_BINARY_CAST_OPERATOR(funcid, name) APPLY_BINARY_OPERATOR(funcid, name)
#define APPLY_TERNARY_CAST_OPERATOR(funcid, name) APPLY_TERNARY_OPERATOR(funcid, name)
llvm::Value* ProcessOpExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitOpExpr* expr, int* max_arg)
static llvm::Value* ProcessOpExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitOpExpr* expr, int* max_arg)
{
llvm::Value* result = nullptr;
@ -835,7 +850,7 @@ llvm::Value* ProcessOpExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitOpEx
return result;
}
llvm::Value* ProcessFuncExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitFuncExpr* expr, int* max_arg)
static llvm::Value* ProcessFuncExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitFuncExpr* expr, int* max_arg)
{
llvm::Value* result = nullptr;
@ -868,12 +883,12 @@ llvm::Value* ProcessFuncExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitFu
#undef APPLY_BINARY_CAST_OPERATOR
#undef APPLY_TERNARY_CAST_OPERATOR
llvm::Value* ProcessSubLinkExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitSubLinkExpr* expr, int* max_arg)
static llvm::Value* ProcessSubLinkExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitSubLinkExpr* expr, int* max_arg)
{
return AddSelectSubQueryResult(ctx, expr->_sub_query_index);
}
llvm::Value* ProcessBoolExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitBoolExpr* expr, int* maxArg)
static llvm::Value* ProcessBoolExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitBoolExpr* expr, int* maxArg)
{
llvm::Value* result = nullptr;
@ -914,7 +929,7 @@ llvm::Value* ProcessBoolExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitBo
return result;
}
llvm::Value* ProcessExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitExpr* expr, int* max_arg)
static llvm::Value* ProcessExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitExpr* expr, int* max_arg)
{
llvm::Value* result = nullptr;
@ -1069,7 +1084,7 @@ bool selectRowColumns(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitSelectExp
return result;
}
bool buildClosedRangeScan(JitLlvmCodeGenContext* ctx, JitIndexScan* index_scan, int* max_arg,
static bool buildClosedRangeScan(JitLlvmCodeGenContext* ctx, JitIndexScan* index_scan, int* max_arg,
JitRangeScanType range_scan_type, llvm::Value* outer_row, int subQueryIndex)
{
// a closed range scan starts just like a point scan (without enough search expressions) and then adds key patterns
@ -1704,7 +1719,7 @@ bool prepareAggregateCount(JitLlvmCodeGenContext* ctx, JitAggregate* aggregate)
return true;
}
bool prepareDistinctSet(JitLlvmCodeGenContext* ctx, const JitAggregate* aggregate)
static bool prepareDistinctSet(JitLlvmCodeGenContext* ctx, const JitAggregate* aggregate)
{
// we need a hash-set according to the aggregated type (preferably but not necessarily linear-probing hash)
// we use an opaque datum type, with a tailor-made hash-function and equals function
@ -1751,7 +1766,7 @@ bool prepareAggregate(JitLlvmCodeGenContext* ctx, JitAggregate* aggregate)
return result;
}
llvm::Value* buildAggregateAvg(
static llvm::Value* buildAggregateAvg(
JitLlvmCodeGenContext* ctx, const JitAggregate* aggregate, llvm::Value* current_aggregate, llvm::Value* var_expr)
{
llvm::Value* aggregate_expr = nullptr;
@ -1807,7 +1822,7 @@ llvm::Value* buildAggregateAvg(
return aggregate_expr;
}
llvm::Value* buildAggregateSum(
static llvm::Value* buildAggregateSum(
JitLlvmCodeGenContext* ctx, const JitAggregate* aggregate, llvm::Value* current_aggregate, llvm::Value* var_expr)
{
llvm::Value* aggregate_expr = nullptr;
@ -1850,7 +1865,7 @@ llvm::Value* buildAggregateSum(
return aggregate_expr;
}
llvm::Value* buildAggregateMax(
static llvm::Value* buildAggregateMax(
JitLlvmCodeGenContext* ctx, const JitAggregate* aggregate, llvm::Value* current_aggregate, llvm::Value* var_expr)
{
llvm::Value* aggregate_expr = nullptr;
@ -1918,7 +1933,7 @@ llvm::Value* buildAggregateMax(
return aggregate_expr;
}
llvm::Value* buildAggregateMin(
static llvm::Value* buildAggregateMin(
JitLlvmCodeGenContext* ctx, const JitAggregate* aggregate, llvm::Value* current_aggregate, llvm::Value* var_expr)
{
llvm::Value* aggregate_expr = nullptr;
@ -1982,7 +1997,7 @@ llvm::Value* buildAggregateMin(
return aggregate_expr;
}
llvm::Value* buildAggregateCount(
static llvm::Value* buildAggregateCount(
JitLlvmCodeGenContext* ctx, const JitAggregate* aggregate, llvm::Value* count_aggregate)
{
llvm::Value* aggregate_expr = nullptr;
@ -2005,7 +2020,7 @@ llvm::Value* buildAggregateCount(
return aggregate_expr;
}
bool buildAggregateMaxMin(JitLlvmCodeGenContext* ctx, JitAggregate* aggregate, llvm::Value* var_expr)
static bool buildAggregateMaxMin(JitLlvmCodeGenContext* ctx, JitAggregate* aggregate, llvm::Value* var_expr)
{
bool result = true;
@ -2038,7 +2053,7 @@ bool buildAggregateMaxMin(JitLlvmCodeGenContext* ctx, JitAggregate* aggregate, l
return result;
}
bool buildAggregateTuple(JitLlvmCodeGenContext* ctx, JitAggregate* aggregate, llvm::Value* var_expr)
static bool buildAggregateTuple(JitLlvmCodeGenContext* ctx, JitAggregate* aggregate, llvm::Value* var_expr)
{
bool result = false;

View File

@ -37,24 +37,9 @@ namespace JitExec {
/** @brief Builds a code segment for checking if soft memory limit has been reached. */
void buildIsSoftMemoryLimitReached(JitLlvmCodeGenContext* ctx);
/** @brief Builds a code segment for writing datum value to a column. */
void buildWriteDatumColumn(JitLlvmCodeGenContext* ctx, llvm::Value* row, int colid, llvm::Value* datum_value);
/** @brief Builds a code segment for writing a row. */
void buildWriteRow(JitLlvmCodeGenContext* ctx, llvm::Value* row, bool isPKey, JitLlvmRuntimeCursor* cursor);
/** @brief Process a join expression (WHERE clause) and generate code to build a search key. */
bool ProcessJoinExpr(JitLlvmCodeGenContext* ctx, Expr* expr, int* column_count, int* column_array, int* max_arg);
/** @brief Process an operator expression (process only "COLUMN equals EXPR" operators). */
bool ProcessJoinOpExpr(
JitLlvmCodeGenContext* ctx, const OpExpr* op_expr, int* column_count, int* column_array, int* max_arg);
/** @brief Process a boolean operator (process only AND operators, since we handle only point queries, or full-prefix
* range update). */
bool ProcessJoinBoolExpr(
JitLlvmCodeGenContext* ctx, const BoolExpr* boolexpr, int* column_count, int* column_array, int* max_arg);
/** @brief Creates a jitted function for code generation. Builds prototype and entry block. */
void CreateJittedFunction(JitLlvmCodeGenContext* ctx, const char* function_name);
@ -71,8 +56,6 @@ llvm::Value* buildCreateNewRow(JitLlvmCodeGenContext* ctx);
llvm::Value* buildSearchRow(
JitLlvmCodeGenContext* ctx, MOT::AccessType access_type, JitRangeScanType range_scan_type, int subQueryIndex = -1);
llvm::Value* buildFilter(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitFilter* filter, int* max_arg);
bool buildFilterRow(
JitLlvmCodeGenContext* ctx, llvm::Value* row, JitFilterArray* filters, int* max_arg, llvm::BasicBlock* next_block);
@ -82,76 +65,11 @@ void buildInsertRow(JitLlvmCodeGenContext* ctx, llvm::Value* row);
/** @brief Adds code to delete a row. */
void buildDeleteRow(JitLlvmCodeGenContext* ctx);
/** @brief Adds code to search for an iterator. */
llvm::Value* buildSearchIterator(JitLlvmCodeGenContext* ctx, JitIndexScanDirection index_scan_direction,
JitRangeBoundMode range_bound_mode, JitRangeScanType range_scan_type, int subQueryIndex = -1);
/** @brief Adds code to search for an iterator. */
llvm::Value* buildBeginIterator(JitLlvmCodeGenContext* ctx, JitRangeScanType rangeScanType, int subQueryIndex = -1);
/** @brief Adds code to get row from iterator. */
llvm::Value* buildGetRowFromIterator(JitLlvmCodeGenContext* ctx, llvm::BasicBlock* endLoopBlock,
MOT::AccessType access_mode, JitIndexScanDirection index_scan_direction, JitLlvmRuntimeCursor* cursor,
JitRangeScanType range_scan_type, int subQueryIndex = -1);
/** @brief Process constant expression. */
llvm::Value* ProcessConstExpr(
JitLlvmCodeGenContext* ctx, const Const* const_value, int& result_type, int arg_pos, int depth, int* max_arg);
/** @brief Process Param expression. */
llvm::Value* ProcessParamExpr(
JitLlvmCodeGenContext* ctx, const Param* param, int& result_type, int arg_pos, int depth, int* max_arg);
/** @brief Process Relabel expression as Param expression. */
llvm::Value* ProcessRelabelExpr(
JitLlvmCodeGenContext* ctx, RelabelType* relabel_type, int& result_type, int arg_pos, int depth, int* max_arg);
/** @brief Proess Var expression. */
llvm::Value* ProcessVarExpr(
JitLlvmCodeGenContext* ctx, const Var* var, int& result_type, int arg_pos, int depth, int* max_arg);
/** @brief Adds call to PG unary operator. */
llvm::Value* AddExecUnaryOperator(
JitLlvmCodeGenContext* ctx, llvm::Value* param, llvm::Constant* unary_operator, int arg_pos);
/** @brief Adds call to PG binary operator. */
llvm::Value* AddExecBinaryOperator(JitLlvmCodeGenContext* ctx, llvm::Value* lhs_param, llvm::Value* rhs_param,
llvm::Constant* binary_operator, int arg_pos);
/** @brief Adds call to PG ternary operator. */
llvm::Value* AddExecTernaryOperator(JitLlvmCodeGenContext* ctx, llvm::Value* param1, llvm::Value* param2,
llvm::Value* param3, llvm::Constant* ternary_operator, int arg_pos);
/** @brief Process operator expression. */
llvm::Value* ProcessOpExpr(
JitLlvmCodeGenContext* ctx, const OpExpr* op_expr, int& result_type, int arg_pos, int depth, int* max_arg);
/** @brief Process function expression. */
llvm::Value* ProcessFuncExpr(
JitLlvmCodeGenContext* ctx, const FuncExpr* func_expr, int& result_type, int arg_pos, int depth, int* max_arg);
llvm::Value* ProcessFilterExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitFilter* filter, int* max_arg);
/** @brief Process an expression. Generates code to evaluate the expression. */
llvm::Value* ProcessExpr(
JitLlvmCodeGenContext* ctx, Expr* expr, int& result_type, int arg_pos, int depth, int* max_arg);
llvm::Value* ProcessConstExpr(JitLlvmCodeGenContext* ctx, const JitConstExpr* expr, int* max_arg);
llvm::Value* ProcessParamExpr(JitLlvmCodeGenContext* ctx, const JitParamExpr* expr, int* max_arg);
llvm::Value* ProcessVarExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, const JitVarExpr* expr, int* max_arg);
llvm::Value* ProcessOpExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitOpExpr* expr, int* max_arg);
llvm::Value* ProcessFuncExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitFuncExpr* expr, int* max_arg);
llvm::Value* ProcessSubLinkExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitSubLinkExpr* expr, int* max_arg);
llvm::Value* ProcessBoolExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitBoolExpr* expr, int* maxArg);
llvm::Value* ProcessExpr(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitExpr* expr, int* max_arg);
bool buildScanExpression(JitLlvmCodeGenContext* ctx, JitColumnExpr* expr, int* max_arg,
JitRangeIteratorType range_itr_type, JitRangeScanType range_scan_type, llvm::Value* outer_row, int subQueryIndex);
@ -164,9 +82,6 @@ bool writeRowColumns(
bool selectRowColumns(JitLlvmCodeGenContext* ctx, llvm::Value* row, JitSelectExprArray* expr_array, int* max_arg,
JitRangeScanType range_scan_type, int subQueryIndex = -1);
bool buildClosedRangeScan(JitLlvmCodeGenContext* ctx, JitIndexScan* index_scan, int* max_arg,
JitRangeScanType range_scan_type, llvm::Value* outer_row, int subQueryIndex);
bool buildSemiOpenRangeScan(JitLlvmCodeGenContext* ctx, JitIndexScan* indexScan, int* maxArg,
JitRangeScanType rangeScanType, JitRangeBoundMode* beginRangeBound, JitRangeBoundMode* endRangeBound,
llvm::Value* outerRow, int subQueryIndex);
@ -201,29 +116,8 @@ bool prepareAggregateMaxMin(JitLlvmCodeGenContext* ctx, JitAggregate* aggregate)
bool prepareAggregateCount(JitLlvmCodeGenContext* ctx, JitAggregate* aggregate);
bool prepareDistinctSet(JitLlvmCodeGenContext* ctx, const JitAggregate* aggregate);
bool prepareAggregate(JitLlvmCodeGenContext* ctx, JitAggregate* aggregate);
llvm::Value* buildAggregateAvg(
JitLlvmCodeGenContext* ctx, const JitAggregate* aggregate, llvm::Value* current_aggregate, llvm::Value* var_expr);
llvm::Value* buildAggregateSum(
JitLlvmCodeGenContext* ctx, const JitAggregate* aggregate, llvm::Value* current_aggregate, llvm::Value* var_expr);
llvm::Value* buildAggregateMax(
JitLlvmCodeGenContext* ctx, const JitAggregate* aggregate, llvm::Value* current_aggregate, llvm::Value* var_expr);
llvm::Value* buildAggregateMin(
JitLlvmCodeGenContext* ctx, const JitAggregate* aggregate, llvm::Value* current_aggregate, llvm::Value* var_expr);
llvm::Value* buildAggregateCount(
JitLlvmCodeGenContext* ctx, const JitAggregate* aggregate, llvm::Value* count_aggregate);
bool buildAggregateMaxMin(JitLlvmCodeGenContext* ctx, JitAggregate* aggregate, llvm::Value* var_expr);
bool buildAggregateTuple(JitLlvmCodeGenContext* ctx, JitAggregate* aggregate, llvm::Value* var_expr);
bool buildAggregateRow(
JitLlvmCodeGenContext* ctx, JitAggregate* aggregate, llvm::Value* row, llvm::BasicBlock* next_block);

View File

@ -22,10 +22,10 @@
* -------------------------------------------------------------------------
*/
#include <algorithm>
// Be sure to include global.h before postgres.h to avoid conflict between libintl.h (included in global.h)
// and c.h (included in postgres.h).
/*
* ATTENTION: Be sure to include global.h before postgres.h to avoid conflict between libintl.h (included in global.h)
* and c.h (included in postgres.h).
*/
#include "global.h"
#include "jit_plan.h"
#include "jit_common.h"
@ -34,6 +34,8 @@
#include "nodes/pg_list.h"
#include "catalog/pg_aggregate.h"
#include <algorithm>
namespace JitExec {
DECLARE_LOGGER(JitPlan, JitExec)

View File

@ -22,8 +22,10 @@
* -------------------------------------------------------------------------
*/
// Be sure to include global.h before postgres.h to avoid conflict between libintl.h (included in global.h)
// and c.h (included in postgres.h).
/*
* ATTENTION: Be sure to include global.h before postgres.h to avoid conflict between libintl.h (included in global.h)
* and c.h (included in postgres.h).
*/
#include "global.h"
#include "jit_plan_expr.h"

View File

@ -31,6 +31,7 @@
#include "storage/mot/jit_def.h"
#include "jit_common.h"
#include "utilities.h"
#include <algorithm>
namespace JitExec {

View File

@ -28,6 +28,7 @@
#include "mot_engine.h"
#include "utilities.h"
#include "mot_error.h"
#include <algorithm>
DECLARE_LOGGER(TVM, JitExec)

View File

@ -22,23 +22,42 @@
* -------------------------------------------------------------------------
*/
// Be sure to include jit_tvm_query.h before anything else because of global.h.
// See jit_tvm_query.h for more details.
#include "jit_tvm_funcs.h"
/*
* ATTENTION: Be sure to include jit_tvm_query.h before anything else because of libintl.h
* (jit_tvm_blocks.h includes jit_tvm_query.h before anything else).
* See jit_tvm_query.h for more details.
*/
#include "jit_tvm_blocks.h"
#include "jit_tvm_funcs.h"
#include "jit_tvm_util.h"
#include "jit_util.h"
#include "mot_error.h"
#include "utilities.h"
#include "catalog/pg_aggregate.h"
using namespace tvm;
namespace JitExec {
DECLARE_LOGGER(JitTvmBlocks, JitExec)
static bool ProcessJoinOpExpr(
JitTvmCodeGenContext* ctx, const OpExpr* op_expr, int* column_count, int* column_array, int* max_arg);
static bool ProcessJoinBoolExpr(
JitTvmCodeGenContext* ctx, const BoolExpr* boolexpr, int* column_count, int* column_array, int* max_arg);
static Instruction* buildExpression(JitTvmCodeGenContext* ctx, Expression* expr);
static Expression* ProcessExpr(
JitTvmCodeGenContext* ctx, Expr* expr, int& result_type, int arg_pos, int depth, int* max_arg);
static Expression* ProcessFilterExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitFilter* filter, int* max_arg);
static Expression* ProcessExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitExpr* expr, int* max_arg);
void CreateJittedFunction(JitTvmCodeGenContext* ctx, const char* function_name, const char* query_string)
{
ctx->m_jittedQuery = ctx->_builder->createFunction(function_name, query_string);
IssueDebugLog("Starting execution of jitted function");
}
bool ProcessJoinExpr(JitTvmCodeGenContext* ctx, Expr* expr, int* column_count, int* column_array, int* max_arg)
static bool ProcessJoinExpr(JitTvmCodeGenContext* ctx, Expr* expr, int* column_count, int* column_array, int* max_arg)
{
bool result = false;
if (expr->type == T_OpExpr) {
@ -51,7 +70,7 @@ bool ProcessJoinExpr(JitTvmCodeGenContext* ctx, Expr* expr, int* column_count, i
return result;
}
bool ProcessJoinOpExpr(
static bool ProcessJoinOpExpr(
JitTvmCodeGenContext* ctx, const OpExpr* op_expr, int* column_count, int* column_array, int* max_arg)
{
bool result = false;
@ -123,7 +142,7 @@ bool ProcessJoinOpExpr(
return result;
}
bool ProcessJoinBoolExpr(
static bool ProcessJoinBoolExpr(
JitTvmCodeGenContext* ctx, const BoolExpr* boolexpr, int* column_count, int* column_array, int* max_arg)
{
bool result = false;
@ -161,7 +180,7 @@ void buildIsSoftMemoryLimitReached(JitTvmCodeGenContext* ctx)
JIT_IF_END()
}
Instruction* buildExpression(JitTvmCodeGenContext* ctx, Expression* expr)
static Instruction* buildExpression(JitTvmCodeGenContext* ctx, Expression* expr)
{
Instruction* expr_value = ctx->_builder->addExpression(expr);
@ -181,7 +200,7 @@ Instruction* buildExpression(JitTvmCodeGenContext* ctx, Expression* expr)
return expr_value;
}
void buildWriteDatumColumn(JitTvmCodeGenContext* ctx, Instruction* row, int colid, Instruction* datum_value)
static void buildWriteDatumColumn(JitTvmCodeGenContext* ctx, Instruction* row, int colid, Instruction* datum_value)
{
// ATTENTION: The datum_value expression-instruction MUST be already evaluated before this code is executed
// That is the reason why the expression-instruction was added to the current block and now we
@ -261,7 +280,7 @@ Instruction* buildSearchRow(JitTvmCodeGenContext* ctx, MOT::AccessType access_ty
return row;
}
Expression* buildFilter(JitTvmCodeGenContext* ctx, Instruction* row, JitFilter* filter, int* max_arg)
static Expression* buildFilter(JitTvmCodeGenContext* ctx, Instruction* row, JitFilter* filter, int* max_arg)
{
Expression* result = nullptr;
Expression* lhs_expr = ProcessExpr(ctx, row, filter->_lhs_operand, max_arg);
@ -339,8 +358,8 @@ void buildDeleteRow(JitTvmCodeGenContext* ctx)
IssueDebugLog("Row deleted");
}
Instruction* buildSearchIterator(JitTvmCodeGenContext* ctx, JitIndexScanDirection index_scan_direction,
JitRangeBoundMode range_bound_mode, JitRangeScanType range_scan_type, int subQueryIndex /* = -1 */)
static Instruction* buildSearchIterator(JitTvmCodeGenContext* ctx, JitIndexScanDirection index_scan_direction,
JitRangeBoundMode range_bound_mode, JitRangeScanType range_scan_type, int subQueryIndex = -1)
{
// search the row
IssueDebugLog("Searching range start");
@ -357,7 +376,8 @@ Instruction* buildSearchIterator(JitTvmCodeGenContext* ctx, JitIndexScanDirectio
}
/** @brief Adds code to search for an iterator. */
Instruction* buildBeginIterator(JitTvmCodeGenContext* ctx, JitRangeScanType rangeScanType, int subQueryIndex /* = -1 */)
static Instruction* buildBeginIterator(
JitTvmCodeGenContext* ctx, JitRangeScanType rangeScanType, int subQueryIndex = -1)
{
// search the row
IssueDebugLog("Getting begin iterator for full-scan");
@ -392,7 +412,7 @@ Instruction* buildGetRowFromIterator(JitTvmCodeGenContext* ctx, BasicBlock* endL
return row;
}
Expression* ProcessConstExpr(
static Expression* ProcessConstExpr(
JitTvmCodeGenContext* ctx, const Const* const_value, int& result_type, int arg_pos, int depth, int* max_arg)
{
Expression* result = nullptr;
@ -416,7 +436,7 @@ Expression* ProcessConstExpr(
return result;
}
Expression* ProcessParamExpr(
static Expression* ProcessParamExpr(
JitTvmCodeGenContext* ctx, const Param* param, int& result_type, int arg_pos, int depth, int* max_arg)
{
Expression* result = nullptr;
@ -440,7 +460,7 @@ Expression* ProcessParamExpr(
return result;
}
Expression* ProcessRelabelExpr(
static Expression* ProcessRelabelExpr(
JitTvmCodeGenContext* ctx, RelabelType* relabel_type, int& result_type, int arg_pos, int depth, int* max_arg)
{
MOT_LOG_DEBUG("Processing RELABEL expression");
@ -448,7 +468,7 @@ Expression* ProcessRelabelExpr(
return ProcessParamExpr(ctx, param, result_type, arg_pos, depth, max_arg);
}
Expression* ProcessVarExpr(
static Expression* ProcessVarExpr(
JitTvmCodeGenContext* ctx, const Var* var, int& result_type, int arg_pos, int depth, int* max_arg)
{
Expression* result = nullptr;
@ -509,7 +529,7 @@ Expression* ProcessVarExpr(
result = new (std::nothrow) name##Operator(args[0], args[1], args[2], arg_pos); \
break;
Expression* ProcessOpExpr(
static Expression* ProcessOpExpr(
JitTvmCodeGenContext* ctx, const OpExpr* op_expr, int& result_type, int arg_pos, int depth, int* max_arg)
{
Expression* result = nullptr;
@ -555,7 +575,7 @@ Expression* ProcessOpExpr(
return result;
}
Expression* ProcessFuncExpr(
static Expression* ProcessFuncExpr(
JitTvmCodeGenContext* ctx, const FuncExpr* func_expr, int& result_type, int arg_pos, int depth, int* max_arg)
{
Expression* result = nullptr;
@ -608,7 +628,8 @@ Expression* ProcessFuncExpr(
#undef APPLY_BINARY_CAST_OPERATOR
#undef APPLY_TERNARY_CAST_OPERATOR
Expression* ProcessExpr(JitTvmCodeGenContext* ctx, Expr* expr, int& result_type, int arg_pos, int depth, int* max_arg)
static Expression* ProcessExpr(
JitTvmCodeGenContext* ctx, Expr* expr, int& result_type, int arg_pos, int depth, int* max_arg)
{
Expression* result = nullptr;
MOT_LOG_DEBUG("%*s --> Processing expression %d", depth, "", (int)expr->type);
@ -639,7 +660,7 @@ Expression* ProcessExpr(JitTvmCodeGenContext* ctx, Expr* expr, int& result_type,
return result;
}
Expression* ProcessConstExpr(JitTvmCodeGenContext* ctx, const JitConstExpr* expr, int* max_arg)
static Expression* ProcessConstExpr(JitTvmCodeGenContext* ctx, const JitConstExpr* expr, int* max_arg)
{
AddSetExprArgIsNull(ctx, expr->_arg_pos, (expr->_is_null ? 1 : 0)); // mark expression null status
Expression* result = new (std::nothrow) ConstExpression(expr->_value, expr->_arg_pos, (int)(expr->_is_null));
@ -649,7 +670,7 @@ Expression* ProcessConstExpr(JitTvmCodeGenContext* ctx, const JitConstExpr* expr
return result;
}
Expression* ProcessParamExpr(JitTvmCodeGenContext* ctx, const JitParamExpr* expr, int* max_arg)
static Expression* ProcessParamExpr(JitTvmCodeGenContext* ctx, const JitParamExpr* expr, int* max_arg)
{
Expression* result = AddGetDatumParam(ctx, expr->_param_id, expr->_arg_pos);
if (max_arg && (expr->_arg_pos > *max_arg)) {
@ -658,7 +679,7 @@ Expression* ProcessParamExpr(JitTvmCodeGenContext* ctx, const JitParamExpr* expr
return result;
}
Expression* ProcessVarExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitVarExpr* expr, int* max_arg)
static Expression* ProcessVarExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitVarExpr* expr, int* max_arg)
{
Expression* result = nullptr;
if (row == nullptr) {
@ -708,7 +729,7 @@ Expression* ProcessVarExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitVarEx
result = new (std::nothrow) name##Operator(args[0], args[1], args[2], arg_pos); \
break;
Expression* ProcessOpExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitOpExpr* expr, int* max_arg)
static Expression* ProcessOpExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitOpExpr* expr, int* max_arg)
{
Expression* result = nullptr;
@ -737,7 +758,7 @@ Expression* ProcessOpExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitOpExpr
return result;
}
Expression* ProcessFuncExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitFuncExpr* expr, int* max_arg)
static Expression* ProcessFuncExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitFuncExpr* expr, int* max_arg)
{
Expression* result = nullptr;
@ -792,7 +813,7 @@ Expression* ProcessFuncExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitFunc
MOT_LOG_TRACE("Unexpected call in filter expression to ternary cast builtin: " #name); \
break;
Expression* ProcessFilterExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitFilter* filter, int* max_arg)
static Expression* ProcessFilterExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitFilter* filter, int* max_arg)
{
Expression* result = nullptr;
@ -830,12 +851,12 @@ Expression* ProcessFilterExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitFi
#undef APPLY_BINARY_CAST_OPERATOR
#undef APPLY_TERNARY_CAST_OPERATOR
Expression* ProcessSubLinkExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitSubLinkExpr* expr, int* max_arg)
static Expression* ProcessSubLinkExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitSubLinkExpr* expr, int* max_arg)
{
return AddSelectSubQueryResult(ctx, expr->_sub_query_index);
}
Expression* ProcessBoolExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitBoolExpr* expr, int* maxArg)
static Expression* ProcessBoolExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitBoolExpr* expr, int* maxArg)
{
Expression* result = nullptr;
@ -871,7 +892,7 @@ Expression* ProcessBoolExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitBool
return result;
}
Expression* ProcessExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitExpr* expr, int* max_arg)
static Expression* ProcessExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitExpr* expr, int* max_arg)
{
Expression* result = nullptr;
@ -897,7 +918,7 @@ Expression* ProcessExpr(JitTvmCodeGenContext* ctx, Instruction* row, JitExpr* ex
return result;
}
bool buildScanExpression(JitTvmCodeGenContext* ctx, JitColumnExpr* expr, int* max_arg,
static bool buildScanExpression(JitTvmCodeGenContext* ctx, JitColumnExpr* expr, int* max_arg,
JitRangeIteratorType range_itr_type, JitRangeScanType range_scan_type, Instruction* outer_row, int subQueryIndex)
{
Expression* value_expr = ProcessExpr(ctx, outer_row, expr->_expr, max_arg);
@ -1025,7 +1046,7 @@ bool selectRowColumns(JitTvmCodeGenContext* ctx, Instruction* row, JitSelectExpr
return true;
}
bool buildClosedRangeScan(JitTvmCodeGenContext* ctx, JitIndexScan* indexScan, int* maxArg,
static bool buildClosedRangeScan(JitTvmCodeGenContext* ctx, JitIndexScan* indexScan, int* maxArg,
JitRangeScanType rangeScanType, Instruction* outerRow, int subQueryIndex)
{
// a closed range scan starts just like a point scan (with not enough search expressions) and then adds key patterns
@ -1079,7 +1100,7 @@ bool buildClosedRangeScan(JitTvmCodeGenContext* ctx, JitIndexScan* indexScan, in
return result;
}
bool buildSemiOpenRangeScan(JitTvmCodeGenContext* ctx, JitIndexScan* index_scan, int* max_arg,
static bool buildSemiOpenRangeScan(JitTvmCodeGenContext* ctx, JitIndexScan* index_scan, int* max_arg,
JitRangeScanType range_scan_type, JitRangeBoundMode* begin_range_bound, JitRangeBoundMode* end_range_bound,
Instruction* outer_row, int subQueryIndex)
{
@ -1205,7 +1226,7 @@ bool buildSemiOpenRangeScan(JitTvmCodeGenContext* ctx, JitIndexScan* index_scan,
return result;
}
bool buildOpenRangeScan(JitTvmCodeGenContext* ctx, JitIndexScan* index_scan, int* max_arg,
static bool buildOpenRangeScan(JitTvmCodeGenContext* ctx, JitIndexScan* index_scan, int* max_arg,
JitRangeScanType range_scan_type, JitRangeBoundMode* begin_range_bound, JitRangeBoundMode* end_range_bound,
Instruction* outer_row, int subQueryIndex)
{
@ -1368,9 +1389,9 @@ bool buildOpenRangeScan(JitTvmCodeGenContext* ctx, JitIndexScan* index_scan, int
return result;
}
bool buildRangeScan(JitTvmCodeGenContext* ctx, JitIndexScan* indexScan, int* maxArg, JitRangeScanType rangeScanType,
JitRangeBoundMode* beginRangeBound, JitRangeBoundMode* endRangeBound, Instruction* outerRow,
int subQueryIndex /* = -1 */)
static bool buildRangeScan(JitTvmCodeGenContext* ctx, JitIndexScan* indexScan, int* maxArg,
JitRangeScanType rangeScanType, JitRangeBoundMode* beginRangeBound, JitRangeBoundMode* endRangeBound,
Instruction* outerRow, int subQueryIndex = -1)
{
bool result = false;
@ -1402,7 +1423,7 @@ bool buildRangeScan(JitTvmCodeGenContext* ctx, JitIndexScan* indexScan, int* max
return result;
}
bool buildPrepareStateScan(JitTvmCodeGenContext* ctx, JitIndexScan* index_scan, int* max_arg,
static bool buildPrepareStateScan(JitTvmCodeGenContext* ctx, JitIndexScan* index_scan, int* max_arg,
JitRangeScanType range_scan_type, Instruction* outer_row)
{
JitRangeBoundMode begin_range_bound = JIT_RANGE_BOUND_NONE;
@ -1445,7 +1466,7 @@ bool buildPrepareStateScan(JitTvmCodeGenContext* ctx, JitIndexScan* index_scan,
return true;
}
bool buildPrepareStateRow(JitTvmCodeGenContext* ctx, MOT::AccessType access_mode, JitIndexScan* index_scan,
static bool buildPrepareStateRow(JitTvmCodeGenContext* ctx, MOT::AccessType access_mode, JitIndexScan* index_scan,
int* max_arg, JitRangeScanType range_scan_type, BasicBlock* next_block)
{
MOT_LOG_DEBUG("Generating select code for stateful range select");
@ -1570,7 +1591,7 @@ JitTvmRuntimeCursor buildRangeCursor(JitTvmCodeGenContext* ctx, JitIndexScan* in
return result;
}
bool prepareAggregateAvg(JitTvmCodeGenContext* ctx, const JitAggregate* aggregate)
static bool prepareAggregateAvg(JitTvmCodeGenContext* ctx, const JitAggregate* aggregate)
{
// although we already have this information in the aggregate descriptor, we still check again
switch (aggregate->_func_id) {
@ -1601,7 +1622,7 @@ bool prepareAggregateAvg(JitTvmCodeGenContext* ctx, const JitAggregate* aggregat
return true;
}
bool prepareAggregateSum(JitTvmCodeGenContext* ctx, const JitAggregate* aggregate)
static bool prepareAggregateSum(JitTvmCodeGenContext* ctx, const JitAggregate* aggregate)
{
switch (aggregate->_func_id) {
case INT8SUMFUNCOID:
@ -1637,13 +1658,13 @@ bool prepareAggregateSum(JitTvmCodeGenContext* ctx, const JitAggregate* aggregat
return true;
}
bool prepareAggregateMaxMin(JitTvmCodeGenContext* ctx, JitAggregate* aggregate)
static bool prepareAggregateMaxMin(JitTvmCodeGenContext* ctx, JitAggregate* aggregate)
{
AddResetAggMaxMinNull(ctx);
return true;
}
bool prepareAggregateCount(JitTvmCodeGenContext* ctx, const JitAggregate* aggregate)
static bool prepareAggregateCount(JitTvmCodeGenContext* ctx, const JitAggregate* aggregate)
{
switch (aggregate->_func_id) {
case 2147: // int8inc_any
@ -1659,7 +1680,7 @@ bool prepareAggregateCount(JitTvmCodeGenContext* ctx, const JitAggregate* aggreg
return true;
}
bool prepareDistinctSet(JitTvmCodeGenContext* ctx, const JitAggregate* aggregate)
static bool prepareDistinctSet(JitTvmCodeGenContext* ctx, const JitAggregate* aggregate)
{
// we need a hash-set according to the aggregated type (preferably but not necessarily linear-probing hash)
// we use an opaque datum type, with a tailor-made hash-function and equals function
@ -1706,7 +1727,7 @@ bool prepareAggregate(JitTvmCodeGenContext* ctx, JitAggregate* aggregate)
return result;
}
Expression* buildAggregateAvg(
static Expression* buildAggregateAvg(
JitTvmCodeGenContext* ctx, const JitAggregate* aggregate, Expression* current_aggregate, Expression* var_expr)
{
Expression* aggregate_expr = nullptr;
@ -1761,7 +1782,7 @@ Expression* buildAggregateAvg(
return aggregate_expr;
}
Expression* buildAggregateSum(
static Expression* buildAggregateSum(
JitTvmCodeGenContext* ctx, const JitAggregate* aggregate, Expression* current_aggregate, Expression* var_expr)
{
Expression* aggregate_expr = nullptr;
@ -1804,7 +1825,7 @@ Expression* buildAggregateSum(
return aggregate_expr;
}
Expression* buildAggregateMax(
static Expression* buildAggregateMax(
JitTvmCodeGenContext* ctx, const JitAggregate* aggregate, Expression* current_aggregate, Expression* var_expr)
{
Expression* aggregate_expr = nullptr;
@ -1872,7 +1893,7 @@ Expression* buildAggregateMax(
return aggregate_expr;
}
Expression* buildAggregateMin(
static Expression* buildAggregateMin(
JitTvmCodeGenContext* ctx, const JitAggregate* aggregate, Expression* current_aggregate, Expression* var_expr)
{
Expression* aggregate_expr = nullptr;
@ -1935,7 +1956,8 @@ Expression* buildAggregateMin(
return aggregate_expr;
}
Expression* buildAggregateCount(JitTvmCodeGenContext* ctx, const JitAggregate* aggregate, Expression* count_aggregate)
static Expression* buildAggregateCount(
JitTvmCodeGenContext* ctx, const JitAggregate* aggregate, Expression* count_aggregate)
{
Expression* aggregate_expr = nullptr;
switch (aggregate->_func_id) {
@ -1957,7 +1979,7 @@ Expression* buildAggregateCount(JitTvmCodeGenContext* ctx, const JitAggregate* a
return aggregate_expr;
}
bool buildAggregateMaxMin(JitTvmCodeGenContext* ctx, JitAggregate* aggregate, Expression* var_expr)
static bool buildAggregateMaxMin(JitTvmCodeGenContext* ctx, JitAggregate* aggregate, Expression* var_expr)
{
bool result = true;
@ -1990,7 +2012,7 @@ bool buildAggregateMaxMin(JitTvmCodeGenContext* ctx, JitAggregate* aggregate, Ex
return result;
}
bool buildAggregateTuple(JitTvmCodeGenContext* ctx, JitAggregate* aggregate, Expression* var_expr)
static bool buildAggregateTuple(JitTvmCodeGenContext* ctx, JitAggregate* aggregate, Expression* var_expr)
{
bool result = false;

View File

@ -25,27 +25,18 @@
#ifndef JIT_TVM_BLOCKS_H
#define JIT_TVM_BLOCKS_H
// Be sure to include jit_tvm_query.h before anything else because of global.h.
// See jit_tvm_query.h for more details.
/*
* ATTENTION: Be sure to include jit_tvm_query.h before anything else because of libintl.h
* See jit_tvm_query.h for more details.
*/
#include "jit_tvm_query.h"
#include "jit_plan.h"
namespace JitExec {
void CreateJittedFunction(JitTvmCodeGenContext* ctx, const char* function_name, const char* query_string);
bool ProcessJoinExpr(JitTvmCodeGenContext* ctx, Expr* expr, int* column_count, int* column_array, int* max_arg);
bool ProcessJoinOpExpr(
JitTvmCodeGenContext* ctx, const OpExpr* op_expr, int* column_count, int* column_array, int* max_arg);
bool ProcessJoinBoolExpr(
JitTvmCodeGenContext* ctx, const BoolExpr* boolexpr, int* column_count, int* column_array, int* max_arg);
void buildIsSoftMemoryLimitReached(JitTvmCodeGenContext* ctx);
tvm::Instruction* buildExpression(JitTvmCodeGenContext* ctx, tvm::Expression* expr);
void buildWriteDatumColumn(JitTvmCodeGenContext* ctx, tvm::Instruction* row, int colid, tvm::Instruction* datum_value);
void buildWriteRow(JitTvmCodeGenContext* ctx, tvm::Instruction* row, bool isPKey, JitTvmRuntimeCursor* cursor);
void buildResetRowsProcessed(JitTvmCodeGenContext* ctx);
@ -57,8 +48,6 @@ tvm::Instruction* buildCreateNewRow(JitTvmCodeGenContext* ctx);
tvm::Instruction* buildSearchRow(
JitTvmCodeGenContext* ctx, MOT::AccessType access_type, JitRangeScanType range_scan_type, int subQueryIndex = -1);
tvm::Expression* buildFilter(JitTvmCodeGenContext* ctx, tvm::Instruction* row, JitFilter* filter, int* max_arg);
bool buildFilterRow(JitTvmCodeGenContext* ctx, tvm::Instruction* row, JitFilterArray* filters, int* max_arg,
tvm::BasicBlock* next_block);
@ -66,60 +55,10 @@ void buildInsertRow(JitTvmCodeGenContext* ctx, tvm::Instruction* row);
void buildDeleteRow(JitTvmCodeGenContext* ctx);
tvm::Instruction* buildSearchIterator(JitTvmCodeGenContext* ctx, JitIndexScanDirection index_scan_direction,
JitRangeBoundMode range_bound_mode, JitRangeScanType range_scan_type, int subQueryIndex = -1);
/** @brief Adds code to search for an iterator. */
tvm::Instruction* buildBeginIterator(JitTvmCodeGenContext* ctx, JitRangeScanType rangeScanType, int subQueryIndex = -1);
tvm::Instruction* buildGetRowFromIterator(JitTvmCodeGenContext* ctx, tvm::BasicBlock* endLoopBlock,
MOT::AccessType access_mode, JitIndexScanDirection index_scan_direction, JitTvmRuntimeCursor* cursor,
JitRangeScanType range_scan_type, int subQueryIndex = -1);
tvm::Expression* ProcessConstExpr(
JitTvmCodeGenContext* ctx, const Const* const_value, int& result_type, int arg_pos, int depth, int* max_arg);
tvm::Expression* ProcessParamExpr(
JitTvmCodeGenContext* ctx, const Param* param, int& result_type, int arg_pos, int depth, int* max_arg);
tvm::Expression* ProcessRelabelExpr(
JitTvmCodeGenContext* ctx, RelabelType* relabel_type, int& result_type, int arg_pos, int depth, int* max_arg);
tvm::Expression* ProcessVarExpr(
JitTvmCodeGenContext* ctx, const Var* var, int& result_type, int arg_pos, int depth, int* max_arg);
tvm::Expression* ProcessOpExpr(
JitTvmCodeGenContext* ctx, const OpExpr* op_expr, int& result_type, int arg_pos, int depth, int* max_arg);
tvm::Expression* ProcessFuncExpr(
JitTvmCodeGenContext* ctx, const FuncExpr* func_expr, int& result_type, int arg_pos, int depth, int* max_arg);
tvm::Expression* ProcessExpr(
JitTvmCodeGenContext* ctx, Expr* expr, int& result_type, int arg_pos, int depth, int* max_arg);
tvm::Expression* ProcessConstExpr(JitTvmCodeGenContext* ctx, const JitConstExpr* expr, int* max_arg);
tvm::Expression* ProcessParamExpr(JitTvmCodeGenContext* ctx, const JitParamExpr* expr, int* max_arg);
tvm::Expression* ProcessVarExpr(JitTvmCodeGenContext* ctx, tvm::Instruction* row, JitVarExpr* expr, int* max_arg);
tvm::Expression* ProcessOpExpr(JitTvmCodeGenContext* ctx, tvm::Instruction* row, JitOpExpr* expr, int* max_arg);
tvm::Expression* ProcessFuncExpr(JitTvmCodeGenContext* ctx, tvm::Instruction* row, JitFuncExpr* expr, int* max_arg);
tvm::Expression* ProcessFilterExpr(JitTvmCodeGenContext* ctx, tvm::Instruction* row, JitFilter* filter, int* max_arg);
tvm::Expression* ProcessSubLinkExpr(
JitTvmCodeGenContext* ctx, tvm::Instruction* row, JitSubLinkExpr* expr, int* max_arg);
tvm::Expression* ProcessBoolExpr(JitTvmCodeGenContext* ctx, tvm::Instruction* row, JitBoolExpr* expr, int* maxArg);
tvm::Expression* ProcessExpr(JitTvmCodeGenContext* ctx, tvm::Instruction* row, JitExpr* expr, int* max_arg);
bool buildScanExpression(JitTvmCodeGenContext* ctx, JitColumnExpr* expr, int* max_arg,
JitRangeIteratorType range_itr_type, JitRangeScanType range_scan_type, tvm::Instruction* outer_row,
int subQueryIndex);
bool buildPointScan(JitTvmCodeGenContext* ctx, JitColumnExprArray* expr_array, int* max_arg,
JitRangeScanType range_scan_type, tvm::Instruction* outer_row, int expr_count = -1, int subQueryIndex = -1);
@ -129,27 +68,6 @@ bool writeRowColumns(
bool selectRowColumns(JitTvmCodeGenContext* ctx, tvm::Instruction* row, JitSelectExprArray* expr_array, int* max_arg,
JitRangeScanType range_scan_type, int subQueryIndex = -1);
bool buildClosedRangeScan(JitTvmCodeGenContext* ctx, JitIndexScan* indexScan, int* maxArg,
JitRangeScanType rangeScanType, tvm::Instruction* outerRow, int subQueryIndex);
bool buildSemiOpenRangeScan(JitTvmCodeGenContext* ctx, JitIndexScan* index_scan, int* max_arg,
JitRangeScanType range_scan_type, JitRangeBoundMode* begin_range_bound, JitRangeBoundMode* end_range_bound,
tvm::Instruction* outer_row, int subQueryIndex);
bool buildOpenRangeScan(JitTvmCodeGenContext* ctx, JitIndexScan* index_scan, int* max_arg,
JitRangeScanType range_scan_type, JitRangeBoundMode* begin_range_bound, JitRangeBoundMode* end_range_bound,
tvm::Instruction* outer_row, int subQueryIndex);
bool buildRangeScan(JitTvmCodeGenContext* ctx, JitIndexScan* indexScan, int* maxArg, JitRangeScanType rangeScanType,
JitRangeBoundMode* beginRangeBound, JitRangeBoundMode* endRangeBound, tvm::Instruction* outerRow,
int subQueryIndex = -1);
bool buildPrepareStateScan(JitTvmCodeGenContext* ctx, JitIndexScan* index_scan, int* max_arg,
JitRangeScanType range_scan_type, tvm::Instruction* outer_row);
bool buildPrepareStateRow(JitTvmCodeGenContext* ctx, MOT::AccessType access_mode, JitIndexScan* index_scan,
int* max_arg, JitRangeScanType range_scan_type, tvm::BasicBlock* next_block);
tvm::Instruction* buildPrepareStateScanRow(JitTvmCodeGenContext* ctx, JitIndexScan* index_scan,
JitRangeScanType range_scan_type, MOT::AccessType access_mode, int* max_arg, tvm::Instruction* outer_row,
tvm::BasicBlock* next_block, tvm::BasicBlock** loop_block);
@ -158,37 +76,8 @@ JitTvmRuntimeCursor buildRangeCursor(JitTvmCodeGenContext* ctx, JitIndexScan* in
JitRangeScanType rangeScanType, JitIndexScanDirection indexScanDirection, tvm::Instruction* outerRow,
int subQueryIndex = -1);
bool prepareAggregateAvg(JitTvmCodeGenContext* ctx, const JitAggregate* aggregate);
bool prepareAggregateSum(JitTvmCodeGenContext* ctx, const JitAggregate* aggregate);
bool prepareAggregateMaxMin(JitTvmCodeGenContext* ctx, JitAggregate* aggregate);
bool prepareAggregateCount(JitTvmCodeGenContext* ctx, const JitAggregate* aggregate);
bool prepareDistinctSet(JitTvmCodeGenContext* ctx, const JitAggregate* aggregate);
bool prepareAggregate(JitTvmCodeGenContext* ctx, JitAggregate* aggregate);
tvm::Expression* buildAggregateAvg(JitTvmCodeGenContext* ctx, const JitAggregate* aggregate,
tvm::Expression* current_aggregate, tvm::Expression* var_expr);
tvm::Expression* buildAggregateSum(JitTvmCodeGenContext* ctx, const JitAggregate* aggregate,
tvm::Expression* current_aggregate, tvm::Expression* var_expr);
tvm::Expression* buildAggregateMax(JitTvmCodeGenContext* ctx, const JitAggregate* aggregate,
tvm::Expression* current_aggregate, tvm::Expression* var_expr);
tvm::Expression* buildAggregateMin(JitTvmCodeGenContext* ctx, const JitAggregate* aggregate,
tvm::Expression* current_aggregate, tvm::Expression* var_expr);
tvm::Expression* buildAggregateCount(
JitTvmCodeGenContext* ctx, const JitAggregate* aggregate, tvm::Expression* count_aggregate);
bool buildAggregateMaxMin(JitTvmCodeGenContext* ctx, JitAggregate* aggregate, tvm::Expression* var_expr);
bool buildAggregateTuple(JitTvmCodeGenContext* ctx, JitAggregate* aggregate, tvm::Expression* var_expr);
bool buildAggregateRow(
JitTvmCodeGenContext* ctx, JitAggregate* aggregate, tvm::Instruction* row, tvm::BasicBlock* next_block);

View File

@ -25,9 +25,12 @@
#ifndef JIT_TVM_FUNCS_H
#define JIT_TVM_FUNCS_H
// Be sure to include jit_tvm_query.h before anything else because of global.h.
// See jit_tvm_query.h for more details.
/*
* ATTENTION: Be sure to include jit_tvm_query.h before anything else because of libintl.h
* See jit_tvm_query.h for more details.
*/
#include "jit_tvm_query.h"
#include "jit_plan_expr.h"
#include "logger.h"
namespace JitExec {

View File

@ -26,42 +26,11 @@
#define JIT_TVM_QUERY_H
/*
* ATTENTION:
* 1. Be sure to include gscodegen.h before anything else to avoid clash with PM definition in datetime.h.
* 2. Be sure to include libintl.h before gscodegen.h to avoid problem with gettext.
* ATTENTION: Be sure to include libintl.h before anything else to avoid problem with gettext.
*/
#include "libintl.h"
#include "codegen/gscodegen.h"
#include "postgres.h"
#include "catalog/pg_operator.h"
#include "utils/fmgroids.h"
#include "nodes/parsenodes.h"
#include "storage/ipc.h"
#include "nodes/pg_list.h"
#include "utils/elog.h"
#include "utils/numeric.h"
#include "utils/numeric_gs.h"
#include "catalog/pg_aggregate.h"
#include "mot_internal.h"
#include "storage/mot/jit_exec.h"
#include "jit_common.h"
#include "jit_tvm.h"
#include "jit_tvm_util.h"
#include "jit_util.h"
#include "jit_plan.h"
#include "mot_engine.h"
#include "utilities.h"
#include "mot_internal.h"
#include "catalog_column_types.h"
#include "mot_error.h"
#include "utilities.h"
#include "mm_session_api.h"
#include <list>
#include <string>
#include <cassert>
namespace JitExec {
/** @struct Holds instructions that evaluate in runtime to begin and end iterators of a cursor. */

View File

@ -22,12 +22,17 @@
* -------------------------------------------------------------------------
*/
// Be sure to include jit_tvm_query.h before anything else because of global.h.
// See jit_tvm_query.h for more details.
/*
* ATTENTION: Be sure to include jit_tvm_query.h before anything else because of libintl.h
* See jit_tvm_query.h for more details.
*/
#include "jit_tvm_query.h"
#include "jit_tvm_query_codegen.h"
#include "jit_tvm_funcs.h"
#include "jit_tvm_blocks.h"
#include "storage/mot/jit_exec.h"
#include "jit_tvm_util.h"
#include "jit_util.h"
using namespace tvm;