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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 14 additions & 0 deletions lib/core/database/mysql_connection.dart
Original file line number Diff line number Diff line change
Expand Up @@ -436,6 +436,20 @@ class MysqlConnection {
await execute(sessionTransactionAccessModeSql(readOnly));
}

/// Session settings of a pool slot: the read-only default and, when the slot
/// has one, a server-side limit on SELECT execution time.
Future<void> configureSession({
required bool readOnly,
Duration? statementTimeout,
}) async {
await setSessionReadOnly(readOnly);
if (statementTimeout != null) {
await execute(
'SET SESSION max_execution_time = ${statementTimeout.inMilliseconds}',
);
}
}

/// `SET SESSION TRANSACTION READ ONLY` / `READ WRITE`.
@visibleForTesting
static String sessionTransactionAccessModeSql(bool readOnly) => readOnly
Expand Down
17 changes: 15 additions & 2 deletions lib/core/database/mysql_connection_pool.dart
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,9 @@ import 'package:querya_desktop/core/database/mysql_connection.dart';
import 'package:querya_desktop/core/storage/local_db.dart';

/// Session policy for pooled connections: browse vs SQL editor (writes).
/// Longest statement an MCP session may run; the server stops it after this.
const Duration mcpStatementTimeout = Duration(seconds: 15);

enum MysqlSessionMode {
/// Tree catalog, stats, and Table Browser SELECT/COUNT.
readOnly,
Expand All @@ -15,10 +18,17 @@ enum MysqlSessionMode {

/// Table Browser Save (own TCP session so START TRANSACTION / SET / USE do not leak).
tableWrite,
/// MCP clients: a read-only session of their own, so an agent never shares
/// the user's session. Statements are bounded by [statementTimeout].
mcp,
}

extension MysqlSessionModeReadOnly on MysqlSessionMode {
bool get isReadOnlySession => this == MysqlSessionMode.readOnly;
bool get isReadOnlySession => this == MysqlSessionMode.readOnly || this == MysqlSessionMode.mcp;

/// Server-side limit for every statement on this session, or null for none.
Duration? get statementTimeout =>
this == MysqlSessionMode.mcp ? mcpStatementTimeout : null;
}

typedef MysqlPoolConnectionFactory = Future<MysqlConnection> Function(
Expand Down Expand Up @@ -144,7 +154,10 @@ class MysqlConnectionPool {
if (inFlight != null) return inFlight;
final attempt = () async {
await entry.connection.connect();
await entry.connection.setSessionReadOnly(mode.isReadOnlySession);
await entry.connection.configureSession(
readOnly: mode.isReadOnlySession,
statementTimeout: mode.statementTimeout,
);
}();
entry.reconnecting = attempt;
return attempt.whenComplete(() {
Expand Down
5 changes: 4 additions & 1 deletion lib/core/database/mysql_service.dart
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,10 @@ Future<MysqlConnection> _defaultCreateAndConnect(
database: database.isEmpty ? null : database,
);
await conn.connect();
await conn.setSessionReadOnly(mode.isReadOnlySession);
await conn.configureSession(
readOnly: mode.isReadOnlySession,
statementTimeout: mode.statementTimeout,
);
return conn;
}

Expand Down
14 changes: 14 additions & 0 deletions lib/core/database/postgres_connection.dart
Original file line number Diff line number Diff line change
Expand Up @@ -389,6 +389,20 @@ class PostgresConnection {
);
}

/// Session settings of a pool slot: the read-only default and, when the slot
/// has one, a server-side statement timeout.
Future<void> configureSession({
required bool readOnly,
Duration? statementTimeout,
}) async {
await setSessionReadOnly(readOnly);
if (statementTimeout != null) {
await execute(
'SET statement_timeout = ${statementTimeout.inMilliseconds}',
);
}
}

/// Tests connectivity and returns a result with an optional error message.
Future<({bool ok, String? error})> testConnection() async {
try {
Expand Down
17 changes: 15 additions & 2 deletions lib/core/database/postgres_connection_pool.dart
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,9 @@ import 'package:querya_desktop/core/database/postgres_connection.dart';
import 'package:querya_desktop/core/storage/local_db.dart';

/// Session policy for pooled connections: browse-only vs ad-hoc SQL (writes).
/// Longest statement an MCP session may run; the server stops it after this.
const Duration mcpStatementTimeout = Duration(seconds: 15);

enum PgSessionMode {
/// `SET default_transaction_read_only = ON` after connect.
/// Tree catalog, stats, and Table Browser SELECT.
Expand All @@ -16,11 +19,18 @@ enum PgSessionMode {

/// Table Browser Save / `REFRESH MATERIALIZED VIEW` (own TCP session).
tableWrite,
/// MCP clients: a read-only session of their own, so an agent never shares
/// the user's session. Statements are bounded by [statementTimeout].
mcp,
}

extension PgSessionModeReadOnly on PgSessionMode {
/// Whether this slot should `SET default_transaction_read_only = ON`.
bool get isReadOnlySession => this == PgSessionMode.readOnly;
bool get isReadOnlySession => this == PgSessionMode.readOnly || this == PgSessionMode.mcp;

/// Server-side limit for every statement on this session, or null for none.
Duration? get statementTimeout =>
this == PgSessionMode.mcp ? mcpStatementTimeout : null;
}

/// Creates a connected [PostgresConnection] for the pool (real or fake in tests).
Expand Down Expand Up @@ -133,7 +143,10 @@ class PostgresConnectionPool {
if (inFlight != null) return inFlight;
final attempt = () async {
await entry.connection.connect();
await entry.connection.setSessionReadOnly(mode.isReadOnlySession);
await entry.connection.configureSession(
readOnly: mode.isReadOnlySession,
statementTimeout: mode.statementTimeout,
);
}();
entry.reconnecting = attempt;
return attempt.whenComplete(() {
Expand Down
5 changes: 4 additions & 1 deletion lib/core/database/postgres_service.dart
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,10 @@ Future<PostgresConnection> _defaultCreateAndConnect(
}) async {
final conn = PostgresConnection.fromConnectionRow(row, database: database);
await conn.connect();
await conn.setSessionReadOnly(mode.isReadOnlySession);
await conn.configureSession(
readOnly: mode.isReadOnlySession,
statementTimeout: mode.statementTimeout,
);
return conn;
}

Expand Down
5 changes: 4 additions & 1 deletion lib/core/database/sqlite_connection_pool.dart
Original file line number Diff line number Diff line change
Expand Up @@ -15,10 +15,13 @@ enum SqliteSessionMode {

/// Table Browser Save (own `Database` so BEGIN/ATTACH do not leak).
tableWrite,
/// MCP clients: a read-only session of their own, so an agent never shares
/// the user's session. Statements are bounded by [statementTimeout].
mcp,
}

extension SqliteSessionModeReadOnly on SqliteSessionMode {
bool get isReadOnlySession => this == SqliteSessionMode.readOnly;
bool get isReadOnlySession => this == SqliteSessionMode.readOnly || this == SqliteSessionMode.mcp;
}

/// Factory to build a connected SQLite connection.
Expand Down
17 changes: 13 additions & 4 deletions lib/core/mcp/mcp_sql_delegates.dart
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,8 @@ import 'package:querya_desktop/features/postgresql/postgres_sql_workspace.dart';
import 'package:querya_desktop/features/sqlite/sqlite_sql_workspace.dart';
import 'package:querya_desktop/features/workspace/sql_execution_delegate.dart';

/// Production [McpDelegateFactory]: the SQL editor's delegates, always on a
/// read-only session (Postgres `default_transaction_read_only`, MySQL
/// Production [McpDelegateFactory]: the SQL editor's delegates, always on the
/// MCP session (its own pool slot) and read-only (Postgres `default_transaction_read_only`, MySQL
/// `SET SESSION TRANSACTION READ ONLY`, SQLite `SQLITE_OPEN_READONLY`).
SqlExecutionDelegate createReadOnlyMcpDelegate(
ConnectionRow row,
Expand All @@ -18,13 +18,22 @@ SqlExecutionDelegate createReadOnlyMcpDelegate(
return PostgresSqlExecutionDelegate(
connectionRow: row,
isReadOnly: true,
isMcp: true,
effectiveDatabaseProvider: () =>
db == null || db.isEmpty ? 'postgres' : db,
autocommitProvider: () => true,
);
case SqlDialect.mysql:
return MysqlSqlExecutionDelegate(connectionRow: row, isReadOnly: true);
return MysqlSqlExecutionDelegate(
connectionRow: row,
isReadOnly: true,
isMcp: true,
);
case SqlDialect.sqlite:
return SqliteSqlExecutionDelegate(connectionRow: row, isReadOnly: true);
return SqliteSqlExecutionDelegate(
connectionRow: row,
isReadOnly: true,
isMcp: true,
);
}
}
12 changes: 10 additions & 2 deletions lib/features/mysql/mysql_sql_workspace.dart
Original file line number Diff line number Diff line change
Expand Up @@ -17,11 +17,19 @@ class MysqlSqlExecutionDelegate extends SqlExecutionDelegate {
MysqlSqlExecutionDelegate({
required this.connectionRow,
required this.isReadOnly,
this.isMcp = false,
});

final ConnectionRow connectionRow;
final bool isReadOnly;

/// MCP delegates run on the MCP session of their own, not the user's.
final bool isMcp;

MysqlSessionMode get _sessionMode => isMcp
? MysqlSessionMode.mcp
: (isReadOnly ? MysqlSessionMode.readOnly : MysqlSessionMode.readWrite);

MysqlLease? _lease;

MysqlLease? get lease => _lease;
Expand All @@ -35,7 +43,7 @@ class MysqlSqlExecutionDelegate extends SqlExecutionDelegate {
final lease = await MysqlService.instance.acquire(
connectionRow,
database: poolDatabaseKey,
mode: isReadOnly ? MysqlSessionMode.readOnly : MysqlSessionMode.readWrite,
mode: _sessionMode,
);
_lease = lease;
}
Expand Down Expand Up @@ -165,7 +173,7 @@ class MysqlSqlExecutionDelegate extends SqlExecutionDelegate {
MysqlService.instance.interrupt(
connectionRow,
database: poolDatabaseKey,
mode: isReadOnly ? MysqlSessionMode.readOnly : MysqlSessionMode.readWrite,
mode: _sessionMode,
);
await _lease?.connection.forceClose();
dropLease();
Expand Down
12 changes: 10 additions & 2 deletions lib/features/postgresql/postgres_sql_workspace.dart
Original file line number Diff line number Diff line change
Expand Up @@ -26,10 +26,18 @@ class PostgresSqlExecutionDelegate extends SqlExecutionDelegate {
required this.isReadOnly,
required this.effectiveDatabaseProvider,
required this.autocommitProvider,
this.isMcp = false,
});

final ConnectionRow connectionRow;
final bool isReadOnly;

/// MCP delegates run on the MCP session of their own, not the user's.
final bool isMcp;

PgSessionMode get _sessionMode => isMcp
? PgSessionMode.mcp
: (isReadOnly ? PgSessionMode.readOnly : PgSessionMode.readWrite);
final String Function() effectiveDatabaseProvider;
final bool Function() autocommitProvider;

Expand All @@ -45,7 +53,7 @@ class PostgresSqlExecutionDelegate extends SqlExecutionDelegate {
final lease = await PostgresService.instance.acquire(
connectionRow,
database: db,
mode: isReadOnly ? PgSessionMode.readOnly : PgSessionMode.readWrite,
mode: _sessionMode,
);
_lease = lease;
_interruptDatabase = db;
Expand Down Expand Up @@ -164,7 +172,7 @@ class PostgresSqlExecutionDelegate extends SqlExecutionDelegate {
PostgresService.instance.interrupt(
connectionRow,
database: _interruptDatabase ?? effectiveDatabaseProvider(),
mode: isReadOnly ? PgSessionMode.readOnly : PgSessionMode.readWrite,
mode: _sessionMode,
);
await _lease?.connection.forceClose();
dropLease();
Expand Down
12 changes: 9 additions & 3 deletions lib/features/sqlite/sqlite_sql_workspace.dart
Original file line number Diff line number Diff line change
Expand Up @@ -17,11 +17,19 @@ class SqliteSqlExecutionDelegate extends SqlExecutionDelegate {
SqliteSqlExecutionDelegate({
required this.connectionRow,
required this.isReadOnly,
this.isMcp = false,
});

final ConnectionRow connectionRow;
final bool isReadOnly;

/// MCP delegates run on the MCP session of their own, not the user's.
final bool isMcp;

SqliteSessionMode get _sessionMode => isMcp
? SqliteSessionMode.mcp
: (isReadOnly ? SqliteSessionMode.readOnly : SqliteSessionMode.readWrite);

SqliteLease? _lease;

SqliteLease? get lease => _lease;
Expand All @@ -32,9 +40,7 @@ class SqliteSqlExecutionDelegate extends SqlExecutionDelegate {
_lease = null;
final lease = await SqliteService.instance.acquire(
connectionRow,
mode: isReadOnly
? SqliteSessionMode.readOnly
: SqliteSessionMode.readWrite,
mode: _sessionMode,
);
_lease = lease;
}
Expand Down
31 changes: 31 additions & 0 deletions test/core/database/mysql_connection_pool_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -266,4 +266,35 @@ void main() {
}
});
});

group('MCP session (#1217)', () {
test('MCP has its own slot, read-only, with a server-side statement timeout',
() async {
expect(MysqlSessionMode.mcp.isReadOnlySession, isTrue);
expect(MysqlSessionMode.readOnly.statementTimeout, isNull);
expect(MysqlSessionMode.mcp.statementTimeout, mcpStatementTimeout);
});

test('an MCP lease never shares the connection of the read-only slot',
() async {
final created = <FakeMysqlConnection>[];
final pool = MysqlConnectionPool(
createAndConnect: (row, {required database, required mode}) async {
final c = FakeMysqlConnection();
await c.connect();
created.add(c);
return c;
},
);
final r = _row();
final ui = await pool.acquire(r,
database: 'app', mode: MysqlSessionMode.readOnly);
final mcp = await pool.acquire(r,
database: 'app', mode: MysqlSessionMode.mcp);
expect(identical(ui.connection, mcp.connection), isFalse);
expect(created.length, 2);
ui.release();
mcp.release();
});
});
}
37 changes: 37 additions & 0 deletions test/core/database/postgres_connection_pool_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -627,4 +627,41 @@ void main() {
}
});
});

group('MCP session (#1217)', () {
test('MCP has its own slot, read-only, with a server-side statement timeout',
() async {
expect(PgSessionMode.mcp.isReadOnlySession, isTrue);
expect(PgSessionMode.readOnly.statementTimeout, isNull);
expect(PgSessionMode.mcp.statementTimeout, mcpStatementTimeout);
final pool = PostgresConnectionPool(
createAndConnect: (row, {required database, required mode}) async =>
FakePostgresConnection(),
);
expect(pool.keyFor(1, 'app', PgSessionMode.mcp),
isNot(pool.keyFor(1, 'app', PgSessionMode.readOnly)));
});

test('an MCP lease never shares the connection of the read-only slot',
() async {
final created = <FakePostgresConnection>[];
final pool = PostgresConnectionPool(
createAndConnect: (row, {required database, required mode}) async {
final c = FakePostgresConnection();
await c.connect();
created.add(c);
return c;
},
);
final r = _row();
final ui = await pool.acquire(r,
database: 'app', mode: PgSessionMode.readOnly);
final mcp = await pool.acquire(r,
database: 'app', mode: PgSessionMode.mcp);
expect(identical(ui.connection, mcp.connection), isFalse);
expect(created.length, 2);
ui.release();
mcp.release();
});
});
}
Loading