fix(export): support PostgreSQL temporary-table result export
This commit is contained in:
parent
2443a550e3
commit
503d6748d4
|
|
@ -7659,7 +7659,7 @@ const gridContextMenuItems = computed<ContextMenuItem[]>(() => {
|
|||
</div>
|
||||
<!-- Truncation warning banner -->
|
||||
<div v-if="showTruncationWarning" class="shrink-0 px-3 py-1 bg-amber-500/10 border-b border-amber-500/20 text-xs text-amber-600 dark:text-amber-400 flex items-center gap-1.5">
|
||||
<span>{{ t("grid.truncatedHint", { count: pageSize }) }}</span>
|
||||
<span>{{ t("grid.truncatedHint", { count: result.rows.length }) }}</span>
|
||||
</div>
|
||||
<!-- Content area: table + side/bottom detail panes -->
|
||||
<div class="flex-1 grid min-h-0 overflow-hidden" :style="contentGridStyle">
|
||||
|
|
@ -8915,7 +8915,7 @@ const gridContextMenuItems = computed<ContextMenuItem[]>(() => {
|
|||
<div v-if="!isErrorResult" class="grid grid-cols-[max-content_minmax(0,1fr)_max-content] items-center gap-2 px-3 py-1 border-t text-xs text-muted-foreground bg-muted/30 shrink-0">
|
||||
<div class="flex min-w-0 items-center gap-2 overflow-hidden">
|
||||
<span v-if="hasData" class="shrink-0">
|
||||
{{ t("grid.totalRows", { count: result.rows.length }) }}
|
||||
{{ t(showTruncationWarning ? "grid.loadedRows" : "grid.totalRows", { count: result.rows.length }) }}
|
||||
<span v-if="typeof displayedTotalRowCount === 'number' && displayedTotalRowCount >= 0" class="text-muted-foreground/70">{{ t("grid.totalRowCount", { count: displayedTotalRowCount }) }}</span>
|
||||
<span v-else-if="totalRowCountBusy" class="text-muted-foreground/70">
|
||||
{{ t("grid.totalRowCountLoading") }}
|
||||
|
|
|
|||
|
|
@ -854,6 +854,7 @@ export default {
|
|||
grid: {
|
||||
rows: "{count} rows",
|
||||
totalRows: "Total {count} rows",
|
||||
loadedRows: "Loaded {count} rows",
|
||||
totalRowCount: "({count} total)",
|
||||
totalRowCountLoading: "(counting...)",
|
||||
loadingMore: "Loading more data...",
|
||||
|
|
@ -1229,7 +1230,7 @@ export default {
|
|||
"metadata-unavailable": "DBX could not load table metadata, so result editing is disabled.",
|
||||
},
|
||||
sortUnsupported: "This SQL does not support full-result sorting. Try again with a single SELECT query.",
|
||||
truncatedHint: "Results truncated to {count} rows. Use the footer pagination or adjust rows per page.",
|
||||
truncatedHint: "Results were truncated after loading {count} rows. Use the footer pagination to browse loaded data; exporting the full result reruns the database query.",
|
||||
},
|
||||
exportProgress: {
|
||||
title: "Exporting Table Data",
|
||||
|
|
|
|||
|
|
@ -803,6 +803,7 @@ export default withEnglishFallback({
|
|||
grid: {
|
||||
rows: "{count} filas",
|
||||
totalRows: "Total {count} filas",
|
||||
loadedRows: "{count} filas cargadas",
|
||||
totalRowCount: "({count} en total)",
|
||||
totalRowCountLoading: "(contando...)",
|
||||
loadingMore: "Cargando más datos...",
|
||||
|
|
@ -1169,7 +1170,7 @@ export default withEnglishFallback({
|
|||
"metadata-unavailable": "DBX no pudo cargar los metadatos de la tabla, por lo que la edición de resultados está deshabilitada.",
|
||||
},
|
||||
sortUnsupported: "Este SQL no admite ordenamiento del resultado completo. Intenta con una consulta SELECT simple.",
|
||||
truncatedHint: "Los resultados están limitados a 10.000 filas. Usa LIMIT/OFFSET en tu consulta para paginar.",
|
||||
truncatedHint: "Los resultados se truncaron después de cargar {count} filas. Usa la paginación inferior para explorar los datos cargados; al exportar el resultado completo se vuelve a consultar la base de datos.",
|
||||
filterBuilderSearchColumns: "Buscar columnas...",
|
||||
filterBuilderNoMatchingColumns: "Sin columnas coincidentes",
|
||||
},
|
||||
|
|
|
|||
|
|
@ -801,6 +801,7 @@ export default withEnglishFallback({
|
|||
grid: {
|
||||
rows: "{count} righe",
|
||||
totalRows: "Totale {count} righe",
|
||||
loadedRows: "{count} righe caricate",
|
||||
totalRowCount: "({count} in totale)",
|
||||
totalRowCountLoading: "(conteggio...)",
|
||||
loadingMore: "Caricamento altri dati...",
|
||||
|
|
@ -1167,7 +1168,7 @@ export default withEnglishFallback({
|
|||
"metadata-unavailable": "DBX could not load table metadata, so result editing is disabled.",
|
||||
},
|
||||
sortUnsupported: "Questo SQL non supporta l'ordinamento sull'intero risultato. Riprova con una query SELECT semplice.",
|
||||
truncatedHint: "Risultati troncati a {count} righe. Usa la paginazione a piè di pagina o regola le righe per pagina.",
|
||||
truncatedHint: "I risultati sono stati troncati dopo il caricamento di {count} righe. Usa la paginazione in basso per consultare i dati caricati; l'esportazione del risultato completo riesegue la query sul database.",
|
||||
filterBuilderSearchColumns: "Cerca campi...",
|
||||
filterBuilderNoMatchingColumns: "Nessun campo corrispondente",
|
||||
},
|
||||
|
|
|
|||
|
|
@ -802,6 +802,7 @@ export default withEnglishFallback({
|
|||
grid: {
|
||||
rows: "{count}行",
|
||||
totalRows: "{count}件表示",
|
||||
loadedRows: "{count}件読み込み済み",
|
||||
totalRowCount: "(全{count}件)",
|
||||
totalRowCountLoading: "(カウント中...)",
|
||||
calculateTotalRows: "総行数をカウント",
|
||||
|
|
@ -1162,7 +1163,7 @@ export default withEnglishFallback({
|
|||
"metadata-unavailable": "DBXがテーブルメタデータを読み込めなかったため、結果の編集は無効です。",
|
||||
},
|
||||
sortUnsupported: "このSQLは完全な結果の並び替えをサポートしていません。単一のSELECTクエリで再試行してください。",
|
||||
truncatedHint: "結果は{count}行に切り詰められました。フッターのページネーションを使用するか、1ページあたりの行数を調整してください。",
|
||||
truncatedHint: "結果は{count}行を読み込んだ時点で切り詰められました。フッターのページネーションで読み込み済みデータを参照できます。完全な結果のエクスポート時はデータベースへ再問い合わせします。",
|
||||
loadingMore: "さらにデータを読み込み中...",
|
||||
allLoaded: "すべて読み込み済み",
|
||||
previewSqlEmpty: "プレビューする保留中のSQL変更はありません",
|
||||
|
|
|
|||
|
|
@ -803,6 +803,7 @@ export default withEnglishFallback({
|
|||
grid: {
|
||||
rows: "{count} linhas",
|
||||
totalRows: "Total de {count} linhas",
|
||||
loadedRows: "{count} linhas carregadas",
|
||||
totalRowCount: "({count} no total)",
|
||||
totalRowCountLoading: "(contando...)",
|
||||
loadingMore: "Carregando mais dados...",
|
||||
|
|
@ -1169,7 +1170,7 @@ export default withEnglishFallback({
|
|||
"metadata-unavailable": "O DBX não conseguiu carregar os metadados da tabela, portanto a edição do resultado está desabilitada.",
|
||||
},
|
||||
sortUnsupported: "Este SQL não suporta a ordenação do resultado completo. Tente novamente com uma única consulta SELECT.",
|
||||
truncatedHint: "Resultados truncados em {count} linhas. Use a paginação no rodapé ou ajuste as linhas por página.",
|
||||
truncatedHint: "Os resultados foram truncados após carregar {count} linhas. Use a paginação no rodapé para navegar pelos dados carregados; exportar o resultado completo executa novamente a consulta no banco de dados.",
|
||||
filterBuilderSearchColumns: "Pesquisar campos...",
|
||||
filterBuilderNoMatchingColumns: "Nenhum campo correspondente",
|
||||
},
|
||||
|
|
|
|||
|
|
@ -856,6 +856,7 @@ export default withEnglishFallback({
|
|||
grid: {
|
||||
rows: "{count} 行",
|
||||
totalRows: "共 {count} 行",
|
||||
loadedRows: "已加载 {count} 行",
|
||||
totalRowCount: "(总计 {count} 行)",
|
||||
totalRowCountLoading: "(统计中...)",
|
||||
loadingMore: "加载更多数据...",
|
||||
|
|
@ -1229,7 +1230,7 @@ export default withEnglishFallback({
|
|||
"metadata-unavailable": "无法读取目标表元数据,暂不能启用结果编辑。",
|
||||
},
|
||||
sortUnsupported: "当前 SQL 不支持全量排序,请改为单条 SELECT 查询后再尝试。",
|
||||
truncatedHint: "结果已截断,仅显示前 {count} 行。可通过底部分页继续加载,或调整每页行数。",
|
||||
truncatedHint: "结果已截断,已加载前 {count} 行。可通过底部分页浏览已加载数据;导出完整结果时会重新查询数据库。",
|
||||
},
|
||||
exportProgress: {
|
||||
title: "导出表数据",
|
||||
|
|
|
|||
|
|
@ -803,6 +803,7 @@ export default withEnglishFallback({
|
|||
grid: {
|
||||
rows: "{count} 列",
|
||||
totalRows: "共 {count} 筆",
|
||||
loadedRows: "已載入 {count} 筆",
|
||||
totalRowCount: "(總計 {count} 筆)",
|
||||
totalRowCountLoading: "(統計中...)",
|
||||
loadingMore: "載入更多資料...",
|
||||
|
|
@ -1169,7 +1170,7 @@ export default withEnglishFallback({
|
|||
"metadata-unavailable": "DBX 無法載入資料表 metadata,因此已停用結果編輯。",
|
||||
},
|
||||
sortUnsupported: "目前 SQL 不支援完整排序,請改為單條 SELECT 查詢後再嘗試。",
|
||||
truncatedHint: "結果已截斷,僅顯示前 {count} 行。可經由底部分頁繼續載入,或調整每頁行數。",
|
||||
truncatedHint: "結果已截斷,已載入前 {count} 筆。可透過底部分頁瀏覽已載入資料;匯出完整結果時會重新查詢資料庫。",
|
||||
filterBuilderSearchColumns: "搜尋欄位...",
|
||||
filterBuilderNoMatchingColumns: "沒有匹配的欄位",
|
||||
},
|
||||
|
|
|
|||
|
|
@ -2502,6 +2502,7 @@ export interface QueryResultExportRequest {
|
|||
schema?: string;
|
||||
sql: string;
|
||||
queryBaseSql: string;
|
||||
setupSql?: string[];
|
||||
databaseType: DatabaseType;
|
||||
useAgentCursor: boolean;
|
||||
filePath: string;
|
||||
|
|
|
|||
|
|
@ -4406,6 +4406,10 @@ export const useQueryStore = defineStore("query", () => {
|
|||
if (!effectiveDbType) return undefined;
|
||||
const useAgentCursor = usesAgentCursorForQuery(conn?.db_type);
|
||||
const queryBaseSql = queryResultBaseSql(tab);
|
||||
const resultStatementIndex = tab.result.statement_index;
|
||||
const batchSql = tab.resultBaseSql ?? tab.lastExecutedSql ?? tab.sql;
|
||||
const batchStatements = effectiveDbType === "postgres" && tab.result.truncated === true && Number.isInteger(resultStatementIndex) && resultStatementIndex! > 0 ? splitSqlStatementRanges(batchSql, effectiveDbType) : [];
|
||||
const setupSql = batchStatements[resultStatementIndex!]?.sql === tab.result.sourceStatement ? batchStatements.slice(0, resultStatementIndex).map((statement) => statement.sql) : undefined;
|
||||
const rowLimit = settings.exportRowLimitEnabled ? settings.exportRowLimit : null;
|
||||
const totalRows = typeof tab.resultTotalRowCount === "number" ? (rowLimit === null ? tab.resultTotalRowCount : Math.min(tab.resultTotalRowCount, rowLimit)) : null;
|
||||
const clientSessionId = tabClientSessionId(tab, "export");
|
||||
|
|
@ -4417,6 +4421,7 @@ export const useQueryStore = defineStore("query", () => {
|
|||
schema: tab.schema,
|
||||
sql,
|
||||
queryBaseSql,
|
||||
setupSql,
|
||||
databaseType: effectiveDbType,
|
||||
useAgentCursor,
|
||||
filePath: options.filePath,
|
||||
|
|
|
|||
|
|
@ -2563,6 +2563,7 @@ pub async fn execute_query_with_max_rows_and_cancel(
|
|||
pub async fn stream_select_query_with_cancel(
|
||||
pool: &Pool,
|
||||
schema: Option<&str>,
|
||||
setup_sql: &[String],
|
||||
sql: &str,
|
||||
max_rows: Option<usize>,
|
||||
cancel_token: Option<CancellationToken>,
|
||||
|
|
@ -2589,28 +2590,72 @@ pub async fn stream_select_query_with_cancel(
|
|||
.await?;
|
||||
}
|
||||
|
||||
let pg_cancel_token = client.cancel_token();
|
||||
let setup_transaction_started = !setup_sql.is_empty();
|
||||
if setup_transaction_started {
|
||||
execute_postgres_infra_statement(&client, "BEGIN", budget.recycle_timeout, "export_setup.begin").await?;
|
||||
}
|
||||
|
||||
let query_timeout = budget.query_timeout;
|
||||
let timeout_error =
|
||||
format!("Query timed out after {} seconds", query_timeout.map_or(0, |timeout| timeout.as_secs()));
|
||||
let progress_clock = Arc::new(StreamProgressClock::new());
|
||||
let progress_clock_for_stream = progress_clock.clone();
|
||||
let mut on_stream_item = |item| {
|
||||
on_item(item)?;
|
||||
progress_clock_for_stream.mark();
|
||||
let setup_result = async {
|
||||
for setup_statement in setup_sql {
|
||||
wait_postgres_query(
|
||||
client.cancel_token(),
|
||||
cancel_context.clone(),
|
||||
cancel_token.clone(),
|
||||
query_timeout,
|
||||
budget.cancel_timeout,
|
||||
async {
|
||||
client.batch_execute(setup_statement).await.map_err(pg_error_to_string)?;
|
||||
Ok(())
|
||||
},
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
Ok(())
|
||||
};
|
||||
let result = await_stream_with_progress_timeout(
|
||||
stream_select_query_inner(&client, sql, row_limit, &mut on_stream_item),
|
||||
query_timeout,
|
||||
progress_clock,
|
||||
cancel_token.as_ref(),
|
||||
timeout_error.clone(),
|
||||
)
|
||||
.await;
|
||||
if result.as_ref().is_err_and(|error| error == &timeout_error || error == crate::query::QUERY_CANCELED) {
|
||||
cancel_postgres_query(pg_cancel_token, cancel_context.as_ref(), budget.cancel_timeout).await;
|
||||
}
|
||||
.await;
|
||||
|
||||
let result = match setup_result {
|
||||
Ok(()) => {
|
||||
let pg_cancel_token = client.cancel_token();
|
||||
let progress_clock = Arc::new(StreamProgressClock::new());
|
||||
let progress_clock_for_stream = progress_clock.clone();
|
||||
let mut on_stream_item = |item| {
|
||||
on_item(item)?;
|
||||
progress_clock_for_stream.mark();
|
||||
Ok(())
|
||||
};
|
||||
let result = await_stream_with_progress_timeout(
|
||||
stream_select_query_inner(&client, sql, row_limit, &mut on_stream_item),
|
||||
query_timeout,
|
||||
progress_clock,
|
||||
cancel_token.as_ref(),
|
||||
timeout_error.clone(),
|
||||
)
|
||||
.await;
|
||||
if result.as_ref().is_err_and(|error| error == &timeout_error || error == crate::query::QUERY_CANCELED) {
|
||||
cancel_postgres_query(pg_cancel_token, cancel_context.as_ref(), budget.cancel_timeout).await;
|
||||
}
|
||||
result
|
||||
}
|
||||
Err(error) => Err(error),
|
||||
};
|
||||
|
||||
let result = if setup_transaction_started {
|
||||
let rollback_result =
|
||||
execute_postgres_infra_statement(&client, "ROLLBACK", budget.cleanup_timeout, "export_setup.rollback")
|
||||
.await;
|
||||
match (result, rollback_result) {
|
||||
(Ok(rows), Ok(_)) => Ok(rows),
|
||||
(Err(query_err), Ok(_)) => Err(query_err),
|
||||
(Ok(_), Err(rollback_err)) => Err(rollback_err),
|
||||
(Err(query_err), Err(rollback_err)) => Err(format!("{query_err}; {rollback_err}")),
|
||||
}
|
||||
} else {
|
||||
result
|
||||
};
|
||||
|
||||
if schema_was_set {
|
||||
let reset_result = reset_postgres_search_path(&client, budget.cleanup_timeout, start).await;
|
||||
|
|
|
|||
|
|
@ -26,8 +26,10 @@ use crate::xlsx_export::{
|
|||
XlsxWorksheetData,
|
||||
};
|
||||
use serde_json::Value;
|
||||
use sqlparser::ast::{GroupByExpr, ObjectNamePart, OrderByKind, SelectItem, SetExpr, Statement, TableFactor};
|
||||
use sqlparser::dialect::GenericDialect;
|
||||
use sqlparser::ast::{
|
||||
GroupByExpr, ObjectName, ObjectNamePart, ObjectType, OrderByKind, SelectItem, SetExpr, Statement, TableFactor,
|
||||
};
|
||||
use sqlparser::dialect::{GenericDialect, PostgreSqlDialect};
|
||||
use sqlparser::parser::Parser;
|
||||
use tokio_util::sync::CancellationToken;
|
||||
|
||||
|
|
@ -61,6 +63,8 @@ pub struct QueryResultExportRequest {
|
|||
pub schema: Option<String>,
|
||||
pub sql: String,
|
||||
pub query_base_sql: String,
|
||||
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
||||
pub setup_sql: Vec<String>,
|
||||
pub database_type: DatabaseType,
|
||||
#[serde(default)]
|
||||
pub use_agent_cursor: bool,
|
||||
|
|
@ -85,6 +89,33 @@ pub struct QueryResultExportRequest {
|
|||
pub date_time_format: Option<String>,
|
||||
}
|
||||
|
||||
fn safe_postgres_temp_setup_sql(setup_sql: &[String]) -> Option<Vec<String>> {
|
||||
if setup_sql.is_empty() {
|
||||
return None;
|
||||
}
|
||||
|
||||
let dialect = PostgreSqlDialect {};
|
||||
let mut temporary_tables: Vec<ObjectName> = Vec::new();
|
||||
for sql in setup_sql {
|
||||
let statements = Parser::parse_sql(&dialect, sql).ok()?;
|
||||
let [statement] = statements.as_slice() else {
|
||||
return None;
|
||||
};
|
||||
match statement {
|
||||
Statement::CreateTable(table) if table.temporary => temporary_tables.push(table.name.clone()),
|
||||
Statement::CreateIndex(index) if temporary_tables.iter().any(|name| name == &index.table_name) => {}
|
||||
Statement::Drop { object_type: ObjectType::Table, names, .. }
|
||||
if !names.is_empty() && names.iter().all(|name| temporary_tables.contains(name)) =>
|
||||
{
|
||||
temporary_tables.retain(|table| !names.contains(table));
|
||||
}
|
||||
_ => return None,
|
||||
}
|
||||
}
|
||||
|
||||
Some(setup_sql.to_vec())
|
||||
}
|
||||
|
||||
fn split_excel_cell_text(value: &str) -> Vec<String> {
|
||||
let mut chunks = Vec::new();
|
||||
let mut current = String::new();
|
||||
|
|
@ -731,9 +762,11 @@ async fn try_export_postgres_query_result_stream(
|
|||
let budget = operation_budget_for_pool_key(state, &pool_key, query_export_timeout(request.timeout_secs)).await;
|
||||
let cancel_context = state.get_postgres_cancel_context(&pool_key).await;
|
||||
|
||||
let setup_sql = safe_postgres_temp_setup_sql(&request.setup_sql).unwrap_or_default();
|
||||
crate::db::postgres::stream_select_query_with_cancel(
|
||||
&pool,
|
||||
request.schema.as_deref(),
|
||||
&setup_sql,
|
||||
&request.sql,
|
||||
stream_row_limit,
|
||||
cancel_token,
|
||||
|
|
@ -1428,6 +1461,35 @@ async fn try_export_sqlserver_query_result_stream(
|
|||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn postgres_temp_setup_accepts_only_session_local_table_operations() {
|
||||
let safe = vec![
|
||||
"CREATE TEMPORARY TABLE t1 AS SELECT 1 AS id".to_string(),
|
||||
"CREATE INDEX t1_id ON t1(id)".to_string(),
|
||||
"CREATE TEMP TABLE t2 AS SELECT id FROM t1".to_string(),
|
||||
"DROP TABLE t1".to_string(),
|
||||
];
|
||||
assert_eq!(safe_postgres_temp_setup_sql(&safe), Some(safe.clone()));
|
||||
|
||||
let parenthesized_ctas = vec![
|
||||
"CREATE TEMPORARY TABLE t1 AS (SELECT CURRENT_DATE AS \u{8d77}\u{4fdd}\u{65e5}\u{671f})".to_string(),
|
||||
"CREATE INDEX t1_1 ON t1(\u{8d77}\u{4fdd}\u{65e5}\u{671f}, \u{7ec8}\u{6b62}\u{65e5}\u{671f})".to_string(),
|
||||
];
|
||||
assert_eq!(safe_postgres_temp_setup_sql(&parenthesized_ctas), Some(parenthesized_ctas.clone()));
|
||||
|
||||
let persistent_create = vec!["CREATE TABLE users_copy AS SELECT * FROM users".to_string()];
|
||||
assert!(safe_postgres_temp_setup_sql(&persistent_create).is_none());
|
||||
|
||||
let persistent_write = vec![
|
||||
"CREATE TEMP TABLE t1 AS SELECT 1 AS id".to_string(),
|
||||
"INSERT INTO audit_log(message) VALUES ('export')".to_string(),
|
||||
];
|
||||
assert!(safe_postgres_temp_setup_sql(&persistent_write).is_none());
|
||||
|
||||
let persistent_index = vec!["CREATE INDEX users_name ON users(name)".to_string()];
|
||||
assert!(safe_postgres_temp_setup_sql(&persistent_index).is_none());
|
||||
}
|
||||
|
||||
fn request(format: &str, row_limit: Option<usize>, total_rows: Option<u64>) -> QueryResultExportRequest {
|
||||
QueryResultExportRequest {
|
||||
export_id: "export-1".to_string(),
|
||||
|
|
@ -1436,6 +1498,7 @@ mod tests {
|
|||
schema: None,
|
||||
sql: "SELECT * FROM users".to_string(),
|
||||
query_base_sql: "SELECT * FROM users".to_string(),
|
||||
setup_sql: Vec::new(),
|
||||
database_type: DatabaseType::Postgres,
|
||||
use_agent_cursor: false,
|
||||
file_path: "out.csv".to_string(),
|
||||
|
|
|
|||
|
|
@ -86,6 +86,7 @@ async fn live_clickhouse_query_result_export_xlsx_streams_random_order_query_onc
|
|||
schema: None,
|
||||
sql: sql.clone(),
|
||||
query_base_sql: sql,
|
||||
setup_sql: Vec::new(),
|
||||
database_type: DatabaseType::ClickHouse,
|
||||
use_agent_cursor: false,
|
||||
file_path: file_path.to_string_lossy().to_string(),
|
||||
|
|
|
|||
|
|
@ -154,6 +154,7 @@ async fn live_mysql_query_result_export_xlsx_streams_single_query_without_duplic
|
|||
schema: None,
|
||||
sql: sql.clone(),
|
||||
query_base_sql: sql,
|
||||
setup_sql: Vec::new(),
|
||||
database_type: DatabaseType::Mysql,
|
||||
use_agent_cursor: false,
|
||||
file_path: file_path.to_string_lossy().to_string(),
|
||||
|
|
@ -238,6 +239,7 @@ async fn live_mysql_xlsx_export_can_outlive_query_timeout_while_rows_keep_arrivi
|
|||
schema: None,
|
||||
sql: sql.clone(),
|
||||
query_base_sql: sql,
|
||||
setup_sql: Vec::new(),
|
||||
database_type: DatabaseType::Mysql,
|
||||
use_agent_cursor: false,
|
||||
file_path: file_path.to_string_lossy().to_string(),
|
||||
|
|
|
|||
|
|
@ -119,6 +119,7 @@ async fn live_postgres_query_result_export_uses_single_streamed_query() {
|
|||
schema: Some(schema.clone()),
|
||||
sql: sql.clone(),
|
||||
query_base_sql: sql,
|
||||
setup_sql: Vec::new(),
|
||||
database_type: DatabaseType::Postgres,
|
||||
use_agent_cursor: false,
|
||||
file_path: file_path.to_string_lossy().to_string(),
|
||||
|
|
@ -155,6 +156,88 @@ async fn live_postgres_query_result_export_uses_single_streamed_query() {
|
|||
assert_eq!(csv.lines().count(), 2051, "unexpected csv row count");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[ignore = "requires DBX_LIVE_POSTGRES_* env vars for temporary-table CSV/XLSX export"]
|
||||
async fn live_postgres_truncated_batch_result_export_replays_safe_temp_setup() {
|
||||
let host = std::env::var("DBX_LIVE_POSTGRES_HOST").unwrap_or_else(|_| "127.0.0.1".to_string());
|
||||
let port = std::env::var("DBX_LIVE_POSTGRES_PORT").ok().and_then(|value| value.parse().ok()).unwrap_or(5432);
|
||||
let user = std::env::var("DBX_LIVE_POSTGRES_USER").unwrap_or_else(|_| "postgres".to_string());
|
||||
let password = std::env::var("DBX_LIVE_POSTGRES_PASSWORD").unwrap_or_default();
|
||||
let database = std::env::var("DBX_LIVE_POSTGRES_DATABASE").unwrap_or_else(|_| "postgres".to_string());
|
||||
let suffix = uuid::Uuid::new_v4().simple().to_string();
|
||||
let short_suffix = &suffix[..8];
|
||||
let connection_id = format!("live-postgres-temp-export-{short_suffix}");
|
||||
let config = live_postgres_config(&connection_id, &host, port, &user, &password, &database);
|
||||
let dir = std::env::temp_dir().join(format!("dbx-live-postgres-temp-export-{short_suffix}"));
|
||||
std::fs::create_dir_all(&dir).unwrap();
|
||||
let storage = Storage::open(&dir.join("storage.db")).await.unwrap();
|
||||
let state = AppState::new(storage);
|
||||
state.configs.write().await.insert(config.id.clone(), config);
|
||||
|
||||
let t1 = format!("dbx_temp_t1_{short_suffix}");
|
||||
let t2 = format!("dbx_temp_t2_{short_suffix}");
|
||||
let t5 = format!("dbx_temp_t5_{short_suffix}");
|
||||
let setup_sql = vec![
|
||||
format!("CREATE TEMPORARY TABLE {t1} AS SELECT i AS id FROM generate_series(1, 30123) AS source(i)"),
|
||||
format!("CREATE INDEX {t1}_id ON {t1}(id)"),
|
||||
format!("CREATE TEMPORARY TABLE {t2} AS SELECT id FROM {t1}"),
|
||||
format!("DROP TABLE {t1}"),
|
||||
format!("CREATE TEMPORARY TABLE {t5} AS SELECT id FROM {t2}"),
|
||||
format!("DROP TABLE {t2}"),
|
||||
];
|
||||
let sql = format!("SELECT id FROM {t5} ORDER BY id");
|
||||
let request = QueryResultExportRequest {
|
||||
export_id: format!("live-postgres-temp-export-csv-{short_suffix}"),
|
||||
connection_id: connection_id.clone(),
|
||||
database: database.clone(),
|
||||
schema: Some("public".to_string()),
|
||||
sql: sql.clone(),
|
||||
query_base_sql: sql.clone(),
|
||||
setup_sql: setup_sql.clone(),
|
||||
database_type: DatabaseType::Postgres,
|
||||
use_agent_cursor: false,
|
||||
file_path: dir.join("result.csv").to_string_lossy().to_string(),
|
||||
format: "csv".to_string(),
|
||||
include_sql_sheet: false,
|
||||
page_size: 2000,
|
||||
row_limit: None,
|
||||
total_rows: None,
|
||||
timeout_secs: Some(30),
|
||||
keyset_optimization_enabled: false,
|
||||
client_session_id: Some(format!("temp-export-csv-{short_suffix}")),
|
||||
execution_id: Some(format!("temp-export-csv-{short_suffix}")),
|
||||
date_time_format: None,
|
||||
};
|
||||
let csv_rows = AtomicU64::new(0);
|
||||
export_query_result_core(&state, &request, None, |progress| {
|
||||
csv_rows.store(progress.rows_exported, Ordering::Relaxed);
|
||||
})
|
||||
.await
|
||||
.expect("export temporary-table result to CSV");
|
||||
let csv = std::fs::read_to_string(&request.file_path).unwrap();
|
||||
assert_eq!(csv_rows.load(Ordering::Relaxed), 30_123);
|
||||
assert_eq!(csv.lines().count(), 30_124);
|
||||
|
||||
let xlsx_path = dir.join("result.xlsx");
|
||||
let xlsx_request = QueryResultExportRequest {
|
||||
export_id: format!("live-postgres-temp-export-xlsx-{short_suffix}"),
|
||||
file_path: xlsx_path.to_string_lossy().to_string(),
|
||||
format: "xlsx".to_string(),
|
||||
client_session_id: Some(format!("temp-export-xlsx-{short_suffix}")),
|
||||
execution_id: Some(format!("temp-export-xlsx-{short_suffix}")),
|
||||
..request
|
||||
};
|
||||
let xlsx_rows = AtomicU64::new(0);
|
||||
export_query_result_core(&state, &xlsx_request, None, |progress| {
|
||||
xlsx_rows.store(progress.rows_exported, Ordering::Relaxed);
|
||||
})
|
||||
.await
|
||||
.expect("export temporary-table result to XLSX");
|
||||
assert_eq!(xlsx_rows.load(Ordering::Relaxed), 30_123);
|
||||
assert!(xlsx_path.metadata().unwrap().len() > 100_000);
|
||||
let _ = std::fs::remove_dir_all(dir);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[ignore = "requires DBX_LIVE_POSTGRES_* env vars for a 650,000-row XLSX export"]
|
||||
async fn live_postgres_xlsx_export_can_outlive_query_timeout_while_rows_keep_arriving() {
|
||||
|
|
@ -181,6 +264,7 @@ async fn live_postgres_xlsx_export_can_outlive_query_timeout_while_rows_keep_arr
|
|||
schema: Some("public".to_string()),
|
||||
sql: sql.to_string(),
|
||||
query_base_sql: sql.to_string(),
|
||||
setup_sql: Vec::new(),
|
||||
database_type: DatabaseType::Postgres,
|
||||
use_agent_cursor: false,
|
||||
file_path: file_path.to_string_lossy().to_string(),
|
||||
|
|
@ -242,6 +326,7 @@ async fn live_postgres_stream_still_times_out_without_progress_and_recovers() {
|
|||
schema: Some("public".to_string()),
|
||||
sql: sql.to_string(),
|
||||
query_base_sql: sql.to_string(),
|
||||
setup_sql: Vec::new(),
|
||||
database_type: DatabaseType::Postgres,
|
||||
use_agent_cursor: false,
|
||||
file_path: file_path.to_string_lossy().to_string(),
|
||||
|
|
|
|||
|
|
@ -435,6 +435,7 @@ async fn live_sqlserver_query_result_export_streams_cte_query_to_csv() {
|
|||
schema: Some("dbo".to_string()),
|
||||
sql: sql.clone(),
|
||||
query_base_sql: sql,
|
||||
setup_sql: Vec::new(),
|
||||
database_type: DatabaseType::SqlServer,
|
||||
use_agent_cursor: false,
|
||||
file_path: file_path.to_string_lossy().to_string(),
|
||||
|
|
|
|||
|
|
@ -95,6 +95,7 @@ async fn live_sqlserver_xlsx_export_can_outlive_query_timeout_while_rows_keep_ar
|
|||
schema: Some("dbo".to_string()),
|
||||
sql: sql.to_string(),
|
||||
query_base_sql: sql.to_string(),
|
||||
setup_sql: Vec::new(),
|
||||
database_type: DatabaseType::SqlServer,
|
||||
use_agent_cursor: false,
|
||||
file_path: file_path.to_string_lossy().to_string(),
|
||||
|
|
|
|||
|
|
@ -4109,6 +4109,48 @@ test("buildQueryResultExportRequest uses exportRowLimit when enabled", async ()
|
|||
}
|
||||
});
|
||||
|
||||
test("buildQueryResultExportRequest includes PostgreSQL batch setup for a truncated temporary-table result", async () => {
|
||||
const restoreStorage = installMemoryStorage();
|
||||
setActivePinia(createPinia());
|
||||
const connectionStore = useConnectionStore();
|
||||
const store = useQueryStore();
|
||||
const originalFetch = globalThis.fetch;
|
||||
|
||||
connectionStore.addEphemeralConnection(conn("conn-1"));
|
||||
const tabId = store.createTab("conn-1", "analytics", "Query", "query", "public");
|
||||
const tab = store.tabs.find((item) => item.id === tabId);
|
||||
assert.ok(tab);
|
||||
const batchSql = ["CREATE TEMPORARY TABLE t1 AS SELECT id FROM events", "CREATE INDEX t1_id ON t1(id)", "SELECT * FROM t1", "DROP TABLE t1"].join(";\n");
|
||||
tab.lastExecutedSql = batchSql;
|
||||
tab.resultBaseSql = batchSql;
|
||||
tab.result = {
|
||||
columns: ["id"],
|
||||
rows: [[1]],
|
||||
affected_rows: 0,
|
||||
execution_time_ms: 1,
|
||||
statement_index: 2,
|
||||
sourceStatement: "SELECT * FROM t1",
|
||||
truncated: true,
|
||||
has_more: false,
|
||||
};
|
||||
|
||||
globalThis.fetch = withConnectionHealthMock(async () => new Response("unexpected request", { status: 500 }));
|
||||
|
||||
try {
|
||||
const request = await store.buildQueryResultExportRequest(tabId, {
|
||||
exportId: "export-temp",
|
||||
filePath: "C:\\tmp\\temp.csv",
|
||||
format: "csv",
|
||||
});
|
||||
|
||||
assert.equal(request?.sql, "SELECT * FROM t1");
|
||||
assert.deepEqual(request?.setupSql, ["CREATE TEMPORARY TABLE t1 AS SELECT id FROM events", "CREATE INDEX t1_id ON t1(id)"]);
|
||||
} finally {
|
||||
globalThis.fetch = originalFetch;
|
||||
restoreStorage();
|
||||
}
|
||||
});
|
||||
|
||||
test("buildQueryResultExportRequest caps progress total when export row limit is enabled", async () => {
|
||||
const restoreStorage = installMemoryStorage();
|
||||
setActivePinia(createPinia());
|
||||
|
|
|
|||
Loading…
Reference in New Issue