From 34dce19c41d95964a0dcfeab06d53645e0742b28 Mon Sep 17 00:00:00 2001 From: he_sk Date: Sun, 27 Sep 2026 04:24:50 +0800 Subject: [PATCH] fix: pin manual transactions to one DM connection --- .github/workflows/integration-tests.yml | 52 +++++++++ docs/api-reference.md | 2 + docs/ci.md | 7 ++ docs/test-results/2026-09-27-ha-routing.md | 15 +++ dpi_bridge/dpi_conn.go | 92 +++++++++++++--- dpi_bridge/dpi_stmt.go | 42 ++++++- .../third_party/chunanyong_dm/PATCHES.md | 11 ++ .../chunanyong_dm/bridge_options.go | 21 ++++ scripts/verify_dm_restart_transaction.py | 103 ++++++++++++++++++ .../integration/test_p1_connection_matrix.py | 64 ++++++++++- 10 files changed, 388 insertions(+), 21 deletions(-) create mode 100644 scripts/verify_dm_restart_transaction.py diff --git a/.github/workflows/integration-tests.yml b/.github/workflows/integration-tests.yml index f3d47a4..c7e8882 100644 --- a/.github/workflows/integration-tests.yml +++ b/.github/workflows/integration-tests.yml @@ -188,3 +188,55 @@ jobs: - name: Show SSL database logs on failure if: failure() run: docker logs --tail 80 dmpython-ci-dm8-ssl || true + + real-dm-restart-arm: + name: Real DM8 restart transaction / ARM Linux / Python 3.10 + runs-on: ubuntu-24.04-arm + timeout-minutes: 20 + steps: + - uses: actions/checkout@v4 + + - uses: actions/setup-go@v5 + with: + go-version: "1.21" + cache-dependency-path: dpi_bridge/go.sum + + - uses: actions/setup-python@v5 + with: + python-version: "3.10" + + - name: Start isolated DM8 and extract build headers + run: | + admin_password="DmPyA1$(openssl rand -hex 12)" + echo "::add-mask::$admin_password" + { + echo "DM_RESTART_PASSWORD=$admin_password" + echo 'DM_RESTART_USER=SYSDBA' + echo 'DM_RESTART_HOST=127.0.0.1' + echo 'DM_RESTART_PORT=15240' + echo 'DM_RESTART_CONTAINER=dmpython-ci-dm8-restart' + } >> "$GITHUB_ENV" + docker pull --platform linux/arm64 \ + yhl452493373/dm8@sha256:5b9d23c04b148d5765d6077d95b64032be64daedabeef9cddc7a4216b9ad0a1a + docker run -d --user root --name dmpython-ci-dm8-restart \ + -e SYSDBA_PWD="$admin_password" -e SYSAUDITOR_PWD="$admin_password" \ + -e CHARSET=1 -e DB_NAME=DMPYRESTART -e INSTANCE_NAME=DMPYRESTART \ + -p 127.0.0.1:15240:5236 \ + yhl452493373/dm8@sha256:5b9d23c04b148d5765d6077d95b64032be64daedabeef9cddc7a4216b9ad0a1a + mkdir -p dpi_include + docker cp dmpython-ci-dm8-restart:/opt/dmdbms/drivers/dpi/include/. dpi_include/ + test -f dpi_include/DPI.h + + - name: Verify restart transaction semantics + timeout-minutes: 10 + run: | + python -m pip install setuptools wheel + if ! python setup.py build_ext --inplace > /tmp/dmpython-restart-build.log 2>&1; then + tail -80 /tmp/dmpython-restart-build.log + exit 1 + fi + PYTHONPATH="$PWD" python scripts/verify_dm_restart_transaction.py + + - name: Show restart database logs on failure + if: failure() + run: docker logs --tail 80 dmpython-ci-dm8-restart || true diff --git a/docs/api-reference.md b/docs/api-reference.md index 12a96ef..ccc6e53 100644 --- a/docs/api-reference.md +++ b/docs/api-reference.md @@ -103,6 +103,8 @@ dmPython.connect( - `shutdown(shutdown_type=dmPython.SHUTDOWN_DEFAULT)` - `explain(statement)` - `ping(reconnect=0)` + +手动提交模式下,同一 Python 连接的语句、`commit()` 和 `rollback()` 固定使用同一条数据库连接。若数据库重启或执行超时使该物理连接失效,提交会报错;应丢弃该连接,重新建连后核对写入结果。自动提交模式中的已打开连接可在一次可见的通信错误后恢复,但失败语句不会被自动重放。 - `__enter__()` - `__exit__(exc_type, exc_value, exc_traceback)` diff --git a/docs/ci.md b/docs/ci.md index 3c7ac9e..8692629 100644 --- a/docs/ci.md +++ b/docs/ci.md @@ -6,6 +6,13 @@ The [data-type and connection-option matrix](test-results/2026-09-26-type-connec A separate ARM Linux job starts DM8 with mandatory SSL and tests a verified encrypted connection on Python 3.10. It also checks paths containing spaces and `&`, rejection of a wrong or missing server certificate pin, and rejection of a plain server when `ssl_path` is requested. The repository does not upload the CI client key as an artifact. +Another isolated ARM job restarts DM8 while a `DECIMAL(30,8)` write is +uncommitted. It requires `commit()` on that manual transaction connection to +fail, verifies that no row was persisted, and checks that an already-open +autocommit connection can resume querying. This runs through +`scripts/verify_dm_restart_transaction.py` without access to the user's Orb +network. + The real-database job also exposes its disposable DM8 container to the BFILE tests. Those tests create a binary file in the container and a database directory with the temporary CI administrator credential, grant the test user read access, then remove both resources. Local runs need `DM_BFILE_TEST_CONTAINER` and `DM_CI_ADMIN_PASSWORD` to run these cases; without them, the BFILE cases are skipped. CI supplies both and treats skips as a gate failure. After all five real-database jobs pass, five macOS ARM jobs build and install wheels for CPython 3.9–3.13. They use the existing `DPI_HEADERS_TAR_B64` repository secret, so fork pull requests run the real-database matrix but skip macOS wheel builds. The required `CI gate` check accepts that documented fork exception; it requires both real-database and wheel jobs for trusted branches. Headers and image archives are not committed or uploaded as artifacts. The standalone `Integration Tests` workflow also runs the full five-version suite nightly and can be started manually. [Five-version results and release rehearsal](test-results/2026-09-25-five-python-release-rehearsal.md) record the exact scope. diff --git a/docs/test-results/2026-09-27-ha-routing.md b/docs/test-results/2026-09-27-ha-routing.md index 5689999..1e3a184 100644 --- a/docs/test-results/2026-09-27-ha-routing.md +++ b/docs/test-results/2026-09-27-ha-routing.md @@ -82,3 +82,18 @@ EP_SELECTION=1 SWITCH_TIMES=1 SWITCH_INTERVAL=0 ``` + +## Existing connection after a single-node restart + +An independent, disposable ARM DM8 instance was restarted while Python 3.10 +held an open autocommit connection. The same Python connection queried again +after one visible communication error. A second experiment inserted an exact +`DECIMAL(30,8)` value in manual-commit mode and verified that a separate +connection could not see it. Before the bridge fix, a database restart made +`commit()` return success even though the row was lost. Manual transactions now +pin one physical connection; `commit()` reports an error after that connection +is lost, and a fresh connection confirms the row was not persisted. The +repeatable check is `scripts/verify_dm_restart_transaction.py`; it also runs +in a dedicated GitHub ARM CI job. Local full real-database regression with a +dedicated test user passed 193 cases. A separate primary/standby takeover with +an already-open connection has still not been verified. diff --git a/dpi_bridge/dpi_conn.go b/dpi_bridge/dpi_conn.go index 1eb53b6..f669a00 100644 --- a/dpi_bridge/dpi_conn.go +++ b/dpi_bridge/dpi_conn.go @@ -44,11 +44,13 @@ import ( // connHandle represents a DPI connection. type connHandle struct { - mu sync.Mutex - env *envHandle - conn *dm.DmConnection // actual Go driver connection - db *sql.DB // holds the sql.DB for lifecycle management - tx driver.Tx // active transaction (nil if none) + mu sync.Mutex + env *envHandle + conn *dm.DmConnection // actual Go driver connection + db *sql.DB // holds the sql.DB for lifecycle management + pinned *sql.Conn // manual transactions must stay on one physical connection + tx driver.Tx // active transaction (nil if none) + modeGeneration uint64 // invalidates statements prepared before an autocommit change // Connection parameters (set before login) host string @@ -117,6 +119,10 @@ func dpi_free_con(hcon C.dhcon) C.DPIRETURN { return DSQL_INVALID_HANDLE } conn.mu.Lock() + if conn.pinned != nil { + conn.pinned.Close() + conn.pinned = nil + } if conn.db != nil { conn.db.Close() conn.db = nil @@ -148,10 +154,33 @@ func dpi_set_con_attr(hcon C.dhcon, attrID C.sdint4, val C.dpointer, valLen C.sd case DSQL_ATTR_LOGIN_PORT: conn.port = intVal case DSQL_ATTR_AUTOCOMMIT: - conn.autocommit = (intVal != 0) - // If already connected, apply autocommit - if conn.conn != nil { - conn.conn.Exec("SET TRANSACTION AUTOCOMMIT "+map[bool]string{true: "ON", false: "OFF"}[conn.autocommit], nil) + enabled := intVal != 0 + changed := enabled != conn.autocommit + if conn.db != nil && changed { + if enabled { + if dbErr := setPinnedAutoCommit(conn.pinned, true); dbErr != nil { + conn.lastErr = diagFromError(dbErr) + return DSQL_ERROR + } + conn.pinned.Close() + conn.pinned = nil + } else { + pinned, dbErr := conn.db.Conn(context.Background()) + if dbErr != nil { + conn.lastErr = diagFromError(dbErr) + return DSQL_ERROR + } + if dbErr = setPinnedAutoCommit(pinned, false); dbErr != nil { + pinned.Close() + conn.lastErr = diagFromError(dbErr) + return DSQL_ERROR + } + conn.pinned = pinned + } + } + conn.autocommit = enabled + if changed { + conn.modeGeneration++ } case DSQL_ATTR_LOGIN_TIMEOUT: if intVal < 0 { @@ -291,6 +320,16 @@ func dpi_set_con_attr(hcon C.dhcon, attrID C.sdint4, val C.dpointer, valLen C.sd return DSQL_SUCCESS } +func setPinnedAutoCommit(pinned *sql.Conn, enabled bool) error { + return pinned.Raw(func(driverConn interface{}) error { + dmConn, ok := driverConn.(*dm.DmConnection) + if !ok { + return fmt.Errorf("unexpected driver connection type: %T", driverConn) + } + return dmConn.SetAutoCommit(enabled) + }) +} + //export dpi_get_con_attr func dpi_get_con_attr(hcon C.dhcon, attrID C.sdint4, val C.dpointer, bufLen C.sdint4, valLen *C.sdint4) C.DPIRETURN { conn, err := getConnHandle(hcon) @@ -397,6 +436,10 @@ func dpi_get_con_attr(hcon C.dhcon, attrID C.sdint4, val C.dpointer, bufLen C.sd dead := C.sdint4(0) // DSQL_CD_FALSE if conn.conn == nil { dead = 1 // DSQL_CD_TRUE + } else if conn.pinned != nil { + if err := conn.pinned.PingContext(context.Background()); err != nil { + dead = 1 + } } else if conn.db != nil { if err := conn.db.Ping(); err != nil { dead = 1 @@ -573,8 +616,8 @@ func dpi_login(hcon C.dhcon, svr *C.sdbyte, user *C.sdbyte, pwd *C.sdbyte) C.DPI } return nil }) - rawConn.Close() if dbErr != nil { + rawConn.Close() db.Close() conn.lastErr = &diagInfo{ errorCode: -1, @@ -583,6 +626,7 @@ func dpi_login(hcon C.dhcon, svr *C.sdbyte, user *C.sdbyte, pwd *C.sdbyte) C.DPI return DSQL_ERROR } if conn.sslPath != "" && dmConn.SSLMode() != 1 { + rawConn.Close() db.Close() conn.lastErr = &diagInfo{errorCode: -1, message: "ssl_path requested, but the server did not negotiate encrypted SSL"} return DSQL_ERROR @@ -590,6 +634,9 @@ func dpi_login(hcon C.dhcon, svr *C.sdbyte, user *C.sdbyte, pwd *C.sdbyte) C.DPI conn.db = db conn.conn = dmConn + if !conn.autocommit { + conn.pinned = rawConn + } if dmConn.CompressionMode() != 0 { conn.compressMsg = 1 } else { @@ -598,19 +645,22 @@ func dpi_login(hcon C.dhcon, svr *C.sdbyte, user *C.sdbyte, pwd *C.sdbyte) C.DPI // Try to get server version var version string - row := db.QueryRow("SELECT BANNER FROM V$VERSION") + row := rawConn.QueryRowContext(context.Background(), "SELECT BANNER FROM V$VERSION") if row.Scan(&version) == nil { conn.serverVersion = version } // Get server encoding var serverCode int32 - row = db.QueryRow("SELECT UNICODE") + row = rawConn.QueryRowContext(context.Background(), "SELECT UNICODE") if row.Scan(&serverCode) == nil { if serverCode == 1 { conn.serverCode = PG_UTF8 } } + if conn.autocommit { + rawConn.Close() + } return DSQL_SUCCESS } @@ -634,6 +684,10 @@ func dpi_logout(hcon C.dhcon) C.DPIRETURN { conn.tx.Rollback() conn.tx = nil } + if conn.pinned != nil { + conn.pinned.Close() + conn.pinned = nil + } if conn.db != nil { conn.db.Close() conn.db = nil @@ -656,7 +710,12 @@ func dpi_commit(hcon C.dhcon) C.DPIRETURN { return DSQL_ERROR } - _, dbErr := conn.db.Exec("COMMIT") + var dbErr error + if conn.pinned != nil { + _, dbErr = conn.pinned.ExecContext(context.Background(), "COMMIT") + } else { + _, dbErr = conn.db.Exec("COMMIT") + } if dbErr != nil { conn.lastErr = &diagInfo{ errorCode: -1, @@ -682,7 +741,12 @@ func dpi_rollback(hcon C.dhcon) C.DPIRETURN { return DSQL_ERROR } - _, dbErr := conn.db.Exec("ROLLBACK") + var dbErr error + if conn.pinned != nil { + _, dbErr = conn.pinned.ExecContext(context.Background(), "ROLLBACK") + } else { + _, dbErr = conn.db.Exec("ROLLBACK") + } if dbErr != nil { conn.lastErr = &diagInfo{ errorCode: -1, diff --git a/dpi_bridge/dpi_stmt.go b/dpi_bridge/dpi_stmt.go index 5be5dae..3e283b2 100644 --- a/dpi_bridge/dpi_stmt.go +++ b/dpi_bridge/dpi_stmt.go @@ -85,8 +85,9 @@ type stmtHandle struct { conn *connHandle // Prepared statement - prepared *sql.Stmt - sql string + prepared *sql.Stmt + sql string + preparedGeneration uint64 // Result set rows *sql.Rows @@ -374,7 +375,7 @@ func dpi_prepare(hstmt C.dhstmt, sqlTxt *C.sdbyte) C.DPIRETURN { return DSQL_ERROR } - prepared, dbErr := stmt.conn.db.Prepare(sqlStr) + prepared, dbErr := prepareOnConnection(stmt.conn, sqlStr) if dbErr != nil { stmt.lastErr = &diagInfo{ errorCode: -1, @@ -383,6 +384,7 @@ func dpi_prepare(hstmt C.dhstmt, sqlTxt *C.sdbyte) C.DPIRETURN { return DSQL_ERROR } stmt.prepared = prepared + stmt.preparedGeneration = stmt.conn.modeGeneration // Count parameters (count '?' in SQL) paramCount := uint16(0) @@ -423,6 +425,17 @@ func dpi_exec(hstmt C.dhstmt) C.DPIRETURN { stmt.lastErr = &diagInfo{errorCode: -1, message: "Statement not prepared"} return DSQL_ERROR } + if stmt.preparedGeneration != stmt.conn.modeGeneration { + stmt.prepared.Close() + stmt.prepared = nil + prepared, dbErr := prepareOnConnection(stmt.conn, stmt.sql) + if dbErr != nil { + stmt.lastErr = diagFromError(dbErr) + return DSQL_ERROR + } + stmt.prepared = prepared + stmt.preparedGeneration = stmt.conn.modeGeneration + } // Build args from parameter bindings args := buildExecArgs(stmt) @@ -458,6 +471,13 @@ func dpi_exec(hstmt C.dhstmt) C.DPIRETURN { return DSQL_SUCCESS } +func prepareOnConnection(conn *connHandle, query string) (*sql.Stmt, error) { + if conn.pinned != nil { + return conn.pinned.PrepareContext(context.Background(), query) + } + return conn.db.Prepare(query) +} + //export dpi_exec_direct func dpi_exec_direct(hstmt C.dhstmt, sqlTxt *C.sdbyte) C.DPIRETURN { stmt, err := getStmtHandle(hstmt) @@ -493,7 +513,13 @@ func dpi_exec_direct(hstmt C.dhstmt, sqlTxt *C.sdbyte) C.DPIRETURN { defer cancel() if isQuery(sqlStr) { - rows, dbErr := stmt.conn.db.QueryContext(ctx, sqlStr) + var rows *sql.Rows + var dbErr error + if stmt.conn.pinned != nil { + rows, dbErr = stmt.conn.pinned.QueryContext(ctx, sqlStr) + } else { + rows, dbErr = stmt.conn.db.QueryContext(ctx, sqlStr) + } if dbErr != nil { stmt.lastErr = diagFromError(dbErr) return DSQL_ERROR @@ -505,7 +531,13 @@ func dpi_exec_direct(hstmt C.dhstmt, sqlTxt *C.sdbyte) C.DPIRETURN { return DSQL_ERROR } } else { - result, dbErr := stmt.conn.db.ExecContext(ctx, sqlStr) + var result sql.Result + var dbErr error + if stmt.conn.pinned != nil { + result, dbErr = stmt.conn.pinned.ExecContext(ctx, sqlStr) + } else { + result, dbErr = stmt.conn.db.ExecContext(ctx, sqlStr) + } if dbErr != nil { stmt.lastErr = diagFromError(dbErr) return DSQL_ERROR diff --git a/dpi_bridge/third_party/chunanyong_dm/PATCHES.md b/dpi_bridge/third_party/chunanyong_dm/PATCHES.md index f4b2c00..ed28249 100644 --- a/dpi_bridge/third_party/chunanyong_dm/PATCHES.md +++ b/dpi_bridge/third_party/chunanyong_dm/PATCHES.md @@ -181,6 +181,17 @@ - Regression: `tests/integration/test_p1_connection_matrix.py` verifies service config paths containing spaces and `&`, and failover after a failed handshake. +## Patch: update autocommit on a live driver connection + +- File: `bridge_options.go` +- The DPI bridge pins one physical connection for manual transactions. Its + runtime autocommit setter updates the driver's wire-protocol flag on that + connection; executing a SQL `SET AUTOCOMMIT` command is not valid through the + driver. Switching back to autocommit first commits a pending transaction. +- Regression: `test_autocommit_toggle_preserves_transaction_boundary` uses two + real DM8 connections. `scripts/verify_dm_restart_transaction.py` proves a + lost manual transaction cannot be reported as committed after restart. + ## Rollback - Remove `replace gitee.com/chunanyong/dm => ./third_party/chunanyong_dm` in `dpi_bridge/go.mod`. diff --git a/dpi_bridge/third_party/chunanyong_dm/bridge_options.go b/dpi_bridge/third_party/chunanyong_dm/bridge_options.go index cd2f171..7ff5f9f 100644 --- a/dpi_bridge/third_party/chunanyong_dm/bridge_options.go +++ b/dpi_bridge/third_party/chunanyong_dm/bridge_options.go @@ -1,5 +1,7 @@ package dm +import "errors" + // CompressionMode reports the negotiated wire-compression mode to the DPI bridge. func (dc *DmConnection) CompressionMode() int { if dc.dmConnector == nil { @@ -12,3 +14,22 @@ func (dc *DmConnection) CompressionMode() int { func (dc *DmConnection) SSLMode() int { return dc.sslEncrypt } + +// SetAutoCommit changes the mode carried on subsequent wire requests. +// Enabling autocommit commits any pending manual transaction first. +func (dc *DmConnection) SetAutoCommit(enabled bool) error { + if dc.dmConnector == nil { + return errors.New("connection is not initialized") + } + if err := dc.checkClosed(); err != nil { + return err + } + if enabled && !dc.autoCommit { + if err := dc.commit(); err != nil { + return err + } + } + dc.dmConnector.autoCommit = enabled + dc.autoCommit = enabled + return nil +} diff --git a/scripts/verify_dm_restart_transaction.py b/scripts/verify_dm_restart_transaction.py new file mode 100644 index 0000000..f0f7d53 --- /dev/null +++ b/scripts/verify_dm_restart_transaction.py @@ -0,0 +1,103 @@ +"""Check that database restart never reports a lost transaction as committed. + +Run only against a disposable DM8 container. The script restarts that container. +""" + +from __future__ import annotations + +import os +import subprocess +import time +import uuid +from decimal import Decimal + +import dmPython + + +AMOUNT = Decimal("12345678901234567890.12345678") + + +def connect(*, autocommit: bool): + return dmPython.connect( + user=os.environ["DM_RESTART_USER"], + password=os.environ["DM_RESTART_PASSWORD"], + server=os.environ.get("DM_RESTART_HOST", "127.0.0.1"), + port=int(os.environ["DM_RESTART_PORT"]), + autoCommit=( + dmPython.DSQL_AUTOCOMMIT_ON + if autocommit + else dmPython.DSQL_AUTOCOMMIT_OFF + ), + login_timeout=3000, + connection_timeout=3, + ) + + +def wait_for_database(): + deadline = time.monotonic() + 120 + while True: + try: + return connect(autocommit=True) + except dmPython.Error: + if time.monotonic() >= deadline: + raise RuntimeError("DM8 did not become ready within 120 seconds") from None + time.sleep(2) + + +def count_rows(conn, table: str) -> int: + with conn.cursor() as cur: + cur.execute(f"SELECT COUNT(*) FROM {table}") + return cur.fetchone()[0] + + +def main(): + container = os.environ["DM_RESTART_CONTAINER"] + table = "DMPY_RESTART_" + uuid.uuid4().hex[:8].upper() + admin = wait_for_database() + try: + with admin.cursor() as cur: + cur.execute(f"CREATE TABLE {table} (ID INT PRIMARY KEY, V DECIMAL(30,8))") + tx = connect(autocommit=False) + try: + with tx.cursor() as cur: + cur.execute(f"INSERT INTO {table} VALUES (?, ?)", (1, AMOUNT)) + assert count_rows(admin, table) == 0, "uncommitted row was visible" + + subprocess.run( + ["docker", "restart", container], + check=True, + stdout=subprocess.DEVNULL, + ) + with wait_for_database() as reader: + assert count_rows(reader, table) == 0, "uncommitted row survived restart" + try: + tx.commit() + except dmPython.Error: + pass + else: + raise AssertionError("commit reported success for a lost transaction") + assert count_rows(reader, table) == 0 + + deadline = time.monotonic() + 30 + errors = 0 + while True: + try: + assert count_rows(admin, table) == 0 + break + except dmPython.Error: + errors += 1 + if time.monotonic() >= deadline: + raise AssertionError("existing autocommit connection did not recover") from None + time.sleep(2) + print(f"lost transaction rejected; existing autocommit connection recovered after {errors} error(s)") + finally: + tx.close() + finally: + with wait_for_database() as cleanup: + with cleanup.cursor() as cur: + cur.execute(f"DROP TABLE {table}") + admin.close() + + +if __name__ == "__main__": + main() diff --git a/tests/integration/test_p1_connection_matrix.py b/tests/integration/test_p1_connection_matrix.py index db29741..8316896 100644 --- a/tests/integration/test_p1_connection_matrix.py +++ b/tests/integration/test_p1_connection_matrix.py @@ -240,6 +240,53 @@ def test_autocommit_on_persists_without_explicit_commit( cur.close() +def test_autocommit_toggle_preserves_transaction_boundary( + conn_params, table_name_factory +): + table = table_name_factory("DMPY_AUTOCOMMIT_TOGGLE") + with dmPython.connect( + **conn_params, autoCommit=dmPython.DSQL_AUTOCOMMIT_ON + ) as writer, dmPython.connect( + **conn_params, autoCommit=dmPython.DSQL_AUTOCOMMIT_ON + ) as reader: + with writer.cursor() as cur: + cur.execute(f"CREATE TABLE {table} (ID INT PRIMARY KEY)") + try: + with writer.cursor() as cur: + cur.prepare(f"INSERT INTO {table} VALUES (?)") + writer.autocommit = dmPython.DSQL_AUTOCOMMIT_OFF + cur.execute(f"INSERT INTO {table} VALUES (?)", (1,)) + with reader.cursor() as cur: + cur.execute(f"SELECT COUNT(*) FROM {table}") + assert cur.fetchone() == (0,) + writer.commit() + with reader.cursor() as cur: + cur.execute(f"SELECT COUNT(*) FROM {table}") + assert cur.fetchone() == (1,) + + writer.autocommit = dmPython.DSQL_AUTOCOMMIT_ON + with writer.cursor() as cur: + cur.execute(f"INSERT INTO {table} VALUES (2)") + with reader.cursor() as cur: + cur.execute(f"SELECT COUNT(*) FROM {table}") + assert cur.fetchone() == (2,) + + with writer.cursor() as cur: + cur.prepare(f"INSERT INTO {table} VALUES (?)") + writer.autocommit = dmPython.DSQL_AUTOCOMMIT_OFF + cur.execute(f"INSERT INTO {table} VALUES (?)", (3,)) + with reader.cursor() as cur: + cur.execute(f"SELECT COUNT(*) FROM {table}") + assert cur.fetchone() == (2,) + writer.autocommit = dmPython.DSQL_AUTOCOMMIT_ON + with reader.cursor() as cur: + cur.execute(f"SELECT COUNT(*) FROM {table}") + assert cur.fetchone() == (3,) + finally: + with reader.cursor() as cur: + cur.execute(f"DROP TABLE {table}") + + def test_isolation_connect_option(conn_params): with dmPython.connect(**conn_params, txn_isolation=dmPython.ISO_LEVEL_READ_COMMITTED) as conn: assert int(conn.txn_isolation) == int(dmPython.ISO_LEVEL_READ_COMMITTED) @@ -462,8 +509,21 @@ def test_connection_timeout_limits_sql_execution( print(time.monotonic() - start) else: raise AssertionError("blocked update unexpectedly completed") -cur.execute("SELECT 1") -assert cur.fetchone() == (1,) +try: + cur.execute("SELECT 1") +except dmPython.Error: + pass +else: + raise AssertionError("lost manual transaction connection stayed usable") +with dmPython.connect( + user=os.environ["DM_TEST_USER"], + password=os.environ["DM_TEST_PASSWORD"], + server=os.environ["DM_TEST_HOST"], + port=int(os.environ["DM_TEST_PORT"]), +) as fresh: + with fresh.cursor() as fresh_cur: + fresh_cur.execute("SELECT 1") + assert fresh_cur.fetchone() == (1,) """ try: cur.execute(f"CREATE TABLE {table} (id INTEGER PRIMARY KEY, v INTEGER)")