Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
27 commits
Select commit Hold shift + click to select a range
357d159
fix(sqlalchemy): infer Decimal128, Int64 and binary column types in r…
aminghadersohi Sep 26, 2026
77877d5
fix(sqlalchemy): round-trip Decimal values through Decimal128
aminghadersohi Sep 26, 2026
fccb8dc
fix: resolve table-qualified column references to the column
aminghadersohi Sep 26, 2026
5a541ec
fix: return correct rows for GROUP BY, IN, LIKE and keyword aliases
aminghadersohi Sep 26, 2026
b14ae06
fix(sqlalchemy): read and bind Uuid columns as BSON UUIDs
aminghadersohi Sep 26, 2026
8ec2dbf
fix(sqlalchemy): return int for BSON int64 in Integer columns
aminghadersohi Sep 26, 2026
2570e33
fix: keep -- inside quoted literals when stripping comments
aminghadersohi Sep 26, 2026
b93eb51
ci: publish immutable wheels from the fork
aminghadersohi Sep 26, 2026
4bdff73
chore: release 0.7.4.1
aminghadersohi Sep 26, 2026
a71b04b
ci: trust the checkout for git introspection during builds
aminghadersohi Sep 26, 2026
d4cb577
ci: mark the checkout safe in the pod's git config
aminghadersohi Sep 26, 2026
2f7dba3
fix: translate NOT and NULL comparisons with SQL three-valued logic
aminghadersohi Sep 26, 2026
9ad3e14
fix: resolve FROM aliases and reject FROM clauses that cannot be tran…
aminghadersohi Sep 26, 2026
b83ab0b
fix(superset): keep booleans and nullable numbers typed in the SQLite…
aminghadersohi Sep 26, 2026
2339f0e
Merge pull request #1 from preset-io/preset/release-0.7.4.1
aminghadersohi Sep 26, 2026
41eda33
fix: return Decimal and int instead of BSON Decimal128 and Int64
aminghadersohi Sep 26, 2026
2e6abd0
fix: translate predicates, aggregates and paging from the parse tree
aminghadersohi Sep 26, 2026
6989e45
chore: release 0.7.4.2
aminghadersohi Sep 26, 2026
c764ab6
fix: bind SET values, LIMIT and OFFSET as parameters; honour LIMIT 0
aminghadersohi Sep 26, 2026
88ca0ef
fix: bind LIKE patterns as parameters instead of inlining them; trans…
aminghadersohi Sep 26, 2026
98b1a0d
Keep %(name)s in inline string literals; LIMIT with literal_binds
aminghadersohi Sep 26, 2026
08d8438
DATE_TRUNC time grains and exact decimals in superset mode
aminghadersohi Sep 26, 2026
da8cc30
List sqlglot with the optional test requirements
aminghadersohi Sep 26, 2026
c296b7b
SQL NULL semantics for SUM and for aggregates over no rows
aminghadersohi Sep 26, 2026
b75d50c
Superset mode: a subquery without rows yields an empty table
aminghadersohi Sep 26, 2026
5bc8d46
ORDER BY an aggregate that is not selected; HAVING next to DATE_TRUNC
aminghadersohi Sep 26, 2026
bbe04c0
Merge pull request #2 from preset-io/fix/parse-tree-predicates
aminghadersohi Sep 26, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
112 changes: 112 additions & 0 deletions Jenkinsfile
Original file line number Diff line number Diff line change
@@ -0,0 +1,112 @@
// Fork publisher. Pull-request wheels are versioned <version>+pr.<number>.<sha>
// (PEP 440-normalized by ci/release_version.py); stable wheels are published only
// from reviewed master. A published artifact is never overwritten.
podTemplate(
imagePullSecrets: ['preset-pull'],
containers: [
containerTemplate(name: 'ci', image: 'preset/ci:latest',
ttyEnabled: true, command: 'cat'),
containerTemplate(name: 'py-ci', image: 'preset/python:3.9.18-2024-02-21-ci',
ttyEnabled: true, command: 'cat'),
// Disposable server for the test suite; reachable on localhost inside the pod.
containerTemplate(name: 'mongo', image: 'mongo:8.0',
envVars: [
envVar(key: 'MONGO_INITDB_ROOT_USERNAME', value: 'admin'),
envVar(key: 'MONGO_INITDB_ROOT_PASSWORD', value: 'secret'),
])
]
) {
node(POD_LABEL) {
checkout scm
def revision = sh(script: 'git rev-parse HEAD', returnStdout: true).trim()
boolean isMaster = env.BRANCH_NAME == 'master'
boolean isPR = env.CHANGE_ID != null
if (!isMaster && !isPR) {
error('Only master and pull-request builds publish; use a PR.')
}

container('py-ci') {
stage('Test and build') {
def args = isMaster ? '' : "${env.CHANGE_ID} ${revision.take(12)}"
sh '''
set -eu
python -m venv .venv
.venv/bin/pip install 'sqlalchemy==2.0.52' 'pymongo==4.17.0' \
'antlr4-python3-runtime==4.13.2' 'jmespath==1.1.0' 'pandas>=2.2,<3' \
'tenacity==9.1.2' 'sqlglot==30.18.0' 'pytest==8.3.5' 'boto3>=1.36,<2' 'packaging==25.0' \
'build==1.4.4' 'setuptools==80.9.0' 'setuptools_scm==8.3.1' 'wheel==0.45.1'
'''
def version = sh(script: ".venv/bin/python ci/release_version.py ${args}",
returnStdout: true).trim()
env.PUBLISH_VERSION = version
env.WHEEL = "pymongosql-${version}-py3-none-any.whl"
env.KEY = "pymongosql/${env.WHEEL}"
sh '''
set -eu
# The checkout is owned by another uid; setuptools_scm runs git during builds.
# The pod is ephemeral, so its global git config is disposable.
git --version
git config --global --add safe.directory "$PWD"
.venv/bin/pip install --no-deps -e .
for attempt in $(seq 1 30); do
.venv/bin/python -c "import pymongo; pymongo.MongoClient('mongodb://admin:secret@localhost:27017', serverSelectionTimeoutMS=2000).admin.command('ping')" && break
sleep 2
done
.venv/bin/python tests/run_test_server.py setup
.venv/bin/python -m pytest -q tests
.venv/bin/pip uninstall -y pymongosql
python - <<'PY'
import os
import re
from pathlib import Path
path = Path('pymongosql/__init__.py')
source = path.read_text()
pattern = re.compile(r'^__version__: str = "[^"]+"$', re.MULTILINE)
assert len(pattern.findall(source)) == 1
path.write_text(pattern.sub('__version__: str = "' + os.environ['PUBLISH_VERSION'] + '"', source))
PY
SOURCE_DATE_EPOCH=$(git -c safe.directory="$PWD" log -1 --format=%ct)
case "$SOURCE_DATE_EPOCH" in
''|*[!0-9]*) echo "Invalid commit timestamp for reproducible build" >&2; exit 1 ;;
esac
export SOURCE_DATE_EPOCH
# Pin the build backend and remove stale output for reproducible retries.
rm -rf build dist pymongosql.egg-info
.venv/bin/python -m build --wheel --no-isolation
test -f "dist/$WHEEL" || { echo "missing dist/$WHEEL"; ls -1 dist; exit 1; }
.venv/bin/python -c "import os, sys, zipfile; names = zipfile.ZipFile('dist/' + os.environ['WHEEL']).namelist(); sys.exit('wheel ships tests or ci' if any(n.startswith(('tests/', 'ci/')) for n in names) else 0)"
.venv/bin/pip install --force-reinstall --no-deps "dist/$WHEEL"
.venv/bin/python - <<'PY'
import importlib.metadata as im
import os
import sqlalchemy as sa
assert im.version('pymongosql') == os.environ['PUBLISH_VERSION']
engine = sa.create_engine('mongodb://user@localhost/db')
assert engine.dialect.name == 'mongodb'
engine.dispose()
PY
sha256sum "dist/$WHEEL"
'''
}
}
container('ci') {
stage('Publish immutable wheel') {
withCredentials([[
$class: 'AmazonWebServicesCredentialsBinding',
credentialsId: 'ci-user',
accessKeyVariable: 'AWS_ACCESS_KEY_ID',
secretKeyVariable: 'AWS_SECRET_ACCESS_KEY'
]]) {
withEnv(["ALLOW_IDENTICAL_PR_ARTIFACT=${isPR && !isMaster}"]) {
sh '''
set -eu
python -m pip install --quiet 'boto3>=1.36,<2'
python ci/publish_wheel.py
'''
}
}
}
}
archiveArtifacts artifacts: 'dist/*.whl,published.sha256', fingerprint: true
}
}
14 changes: 14 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -743,6 +743,20 @@ PyMongoSQL can be used as a database driver in Apache Superset for querying and

This allows seamless integration between MongoDB data and Superset's BI capabilities without requiring data migration to traditional SQL databases.

**Time grains and decimals:**

- `DATE_TRUNC('<unit>', field)` is translated to MongoDB's `$dateTrunc` (MongoDB 5.0+), in
projections and `GROUP BY`. Units: `second`, `minute`, `hour`, `day`, `week` (starting
Sunday), `week_monday`, `month`, `quarter`, `year`, `week_ending_saturday` and
`week_ending_sunday`. Truncation is in UTC. The same function is available in the
superset-mode SQLite stage, so virtual datasets group by time the same way.
- In superset mode, a subquery's result is loaded into an in-memory SQLite database. Columns
holding `Decimal128` values are evaluated exactly there (install `pymongosql[superset]`,
which adds `sqlglot`): the column itself, `SUM`/`AVG`/`MIN`/`MAX` over it, `GROUP BY`,
`ORDER BY` and comparisons with numeric literals, with Decimal128's 34 significant
digits. Any other use of such a column (arithmetic, other functions, `DISTINCT`
aggregates) raises `NotSupportedError` rather than computing with doubles.

**Important Note on Collection Names:**

When using collection names containing special characters (`.`, `-`, `:`), you must wrap them in double quotes to prevent Superset's SQL parser from incorrectly interpreting them.
Expand Down
71 changes: 71 additions & 0 deletions ci/publish_wheel.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
"""Publish an immutable wheel; retries may reuse identical archive content."""

import hashlib
import io
import os
from pathlib import Path
from zipfile import BadZipFile, ZipFile

import boto3
from botocore.exceptions import ClientError


def same_wheel_content(stored, fresh):
"""Compare sorted member names and bytes, ignoring ZIP metadata such as timestamps."""
if stored == fresh:
return True
try:
with ZipFile(io.BytesIO(stored)) as old, ZipFile(io.BytesIO(fresh)) as new:
old_members = sorted(old.infolist(), key=lambda member: member.filename)
new_members = sorted(new.infolist(), key=lambda member: member.filename)
return [member.filename for member in old_members] == [member.filename for member in new_members] and all(
old.read(a) == new.read(b) for a, b in zip(old_members, new_members)
)
except BadZipFile:
return False


def publish_wheel(s3, bucket, key, body, is_pr=False):
try:
s3.put_object(Bucket=bucket, Key=key, Body=body, IfNoneMatch="*")
except ClientError as error:
if error.response["Error"]["Code"] != "PreconditionFailed":
raise
print("Artifact already exists; verifying archive content without overwriting.")

response = s3.get_object(Bucket=bucket, Key=key)
try:
stored = response["Body"].read()
finally:
response["Body"].close()
if not same_wheel_content(stored, body):
raise RuntimeError(
"Stored wheel differs from this commit's freshly built artifact; "
"refusing to overwrite or accept it (local sha256={}, stored sha256={}).".format(
hashlib.sha256(body).hexdigest(),
hashlib.sha256(stored).hexdigest(),
)
+ ("" if is_pr else " Bump __version__ before publishing different content.")
)
print("Published wheel verified by archive content: " + key)
return hashlib.sha256(stored).hexdigest()


def main():
wheel = os.environ["WHEEL"]
receipt = Path("published.sha256")
# Do not leave a previous build's success receipt after a failed retry.
if receipt.exists():
receipt.unlink()
digest = publish_wheel(
boto3.client("s3"),
"preset-pypi",
os.environ["KEY"],
Path("dist", wheel).read_bytes(),
is_pr=os.environ.get("ALLOW_IDENTICAL_PR_ARTIFACT") == "true",
)
receipt.write_text(digest + " " + wheel + "\n")


if __name__ == "__main__":
main()
39 changes: 39 additions & 0 deletions ci/release_version.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
"""Compute the published version of this fork, normalized as the wheel will be.

Stable builds publish the declared ``__version__`` (a four-part Preset release
such as 0.7.4.1). Pull-request builds publish ``<version>+pr.<change>.<sha>``.
PEP 440 normalizes local segments (lower case; a numeric segment loses leading
zeros), so the filename must be derived from the normalized form or it will not
match the file the build backend writes.
"""

import re
import sys
from pathlib import Path

from packaging.version import Version

INIT = Path(__file__).resolve().parents[1] / "pymongosql" / "__init__.py"


def declared_version(source=None):
source = INIT.read_text() if source is None else source
matches = re.findall(r'^__version__: str = "([^"]+)"$', source, re.MULTILINE)
if len(matches) != 1:
raise SystemExit("expected exactly one __version__ declaration")
version = Version(matches[0])
if str(version) != matches[0] or version.local or len(version.release) != 4:
raise SystemExit(f"__version__ must be a normalized four-part release, got {matches[0]!r}")
return matches[0]


def release_version(base, change_id=None, revision=None):
if change_id is None:
return str(Version(base))
if not re.fullmatch(r"[0-9]+", change_id) or not re.fullmatch(r"[0-9a-fA-F]{7,40}", revision or ""):
raise SystemExit("pull-request builds need a numeric change id and a git revision")
return str(Version(f"{base}+pr.{change_id}.{revision}"))


if __name__ == "__main__":
print(release_version(declared_version(), *sys.argv[1:]))
2 changes: 1 addition & 1 deletion pymongosql/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@
if TYPE_CHECKING:
from .connection import Connection

__version__: str = "0.7.4"
__version__: str = "0.7.4.2"

# Globals https://www.python.org/dev/peps/pep-0249/#globals
apilevel: str = "2.0"
Expand Down
73 changes: 47 additions & 26 deletions pymongosql/executor.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,14 @@ def _run_db_command(db: Any, command: Dict[str, Any], connection: Any, operation
)


def _paging(limit: Any, skip: Any) -> Any:
"""Validate bound LIMIT/OFFSET values; returns (limit, skip)."""
for name, value in (("LIMIT", limit), ("OFFSET", skip)):
if value is not None and (isinstance(value, bool) or not isinstance(value, int) or value < 0):
raise ProgrammingError(f"{name} must be a non-negative integer, got {value!r}")
return limit, skip


@dataclass
class ExecutionContext:
"""Manages execution context for a single query"""
Expand Down Expand Up @@ -149,9 +157,16 @@ def _execute_find_plan(
# Replace placeholders with parameters in filter_stage only (not in projection)
filter_stage = execution_plan.filter_stage or {}

if parameters:
# Positional parameters with ? (named parameters are converted to positional in execute())
filter_stage = self._replace_placeholders(filter_stage, parameters)
# Positional parameters (named ones are converted to positional in execute()),
# in statement order: WHERE, then LIMIT, then OFFSET
bound, _ = SQLHelper.bind_filter(
{"filter": filter_stage, "limit": execution_plan.limit_stage, "skip": execution_plan.skip_stage},
parameters,
)
filter_stage = bound["filter"]
limit, skip = _paging(bound["limit"], bound["skip"])
if limit == 0:
return {"cursor": {"id": 0, "firstBatch": []}, "ok": 1}

projection_stage = execution_plan.projection_stage or {}

Expand All @@ -170,13 +185,11 @@ def _execute_find_plan(
sort_spec[field_name] = direction
find_command["sort"] = sort_spec

# Apply skip if specified
if execution_plan.skip_stage:
find_command["skip"] = execution_plan.skip_stage

# Apply limit if specified
if execution_plan.limit_stage:
find_command["limit"] = execution_plan.limit_stage
# Apply skip and limit if specified (MongoDB reads limit 0 as "no limit")
if skip:
find_command["skip"] = skip
if limit is not None:
find_command["limit"] = limit

_logger.debug(f"Executing MongoDB command: {find_command}")

Expand Down Expand Up @@ -223,7 +236,12 @@ def _execute_aggregate_plan(

# Parse pipeline and options from JSON strings
try:
pipeline = json.loads(execution_plan.aggregate_pipeline or "[]")
if execution_plan.aggregate_parameterized:
from bson import json_util

pipeline = json_util.loads(execution_plan.aggregate_pipeline or "[]")
else:
pipeline = json.loads(execution_plan.aggregate_pipeline or "[]")
options = json.loads(execution_plan.aggregate_options or "{}")
except json.JSONDecodeError as e:
raise ProgrammingError(f"Invalid JSON in aggregate pipeline or options: {e}")
Expand All @@ -232,6 +250,14 @@ def _execute_aggregate_plan(
_logger.debug(f"Pipeline: {pipeline}")
_logger.debug(f"Options: {options}")

# A pipeline generated from SQL carries parameter markers (WHERE, HAVING), then
# LIMIT and OFFSET
limit, skip = execution_plan.limit_stage, execution_plan.skip_stage
if execution_plan.aggregate_parameterized:
bound, _ = SQLHelper.bind_filter({"pipeline": pipeline, "limit": limit, "skip": skip}, parameters)
pipeline, limit, skip = bound["pipeline"], bound["limit"], bound["skip"]
limit, skip = _paging(limit, skip)

# Get collection and call aggregate()
collection = db[execution_plan.collection]

Expand All @@ -258,11 +284,11 @@ def _execute_aggregate_plan(
results = sorted(results, key=lambda x: x.get(field_name), reverse=reverse)

# Apply skip and limit
if execution_plan.skip_stage:
results = results[execution_plan.skip_stage :]
if skip:
results = results[skip:]

if execution_plan.limit_stage:
results = results[: execution_plan.limit_stage]
if limit is not None:
results = results[:limit]

# Apply projection if specified
if execution_plan.projection_stage:
Expand Down Expand Up @@ -513,11 +539,8 @@ def _execute_execution_plan(

filter_conditions = execution_plan.filter_conditions or {}

# Replace placeholders in filter if parameters provided
if parameters and filter_conditions:
filter_conditions = SQLHelper.replace_placeholders_generic(
filter_conditions, parameters, execution_plan.parameter_style
)
# Bind the WHERE clause's parameter markers; every parameter must be used
filter_conditions, _ = SQLHelper.bind_filter(filter_conditions, parameters)

command = {"delete": execution_plan.collection, "deletes": [{"q": filter_conditions, "limit": 0}]}

Expand Down Expand Up @@ -600,12 +623,10 @@ def _execute_execution_plan(
# Replace placeholders if parameters provided
# Note: We need to replace both update_fields and filter_conditions in one pass
# to maintain correct parameter ordering (SET clause first, then WHERE clause)
if parameters:
# Combine structures for replacement in correct order
combined = {"update_fields": update_fields, "filter_conditions": filter_conditions}
replaced = SQLHelper.replace_placeholders_generic(combined, parameters, execution_plan.parameter_style)
update_fields = replaced["update_fields"]
filter_conditions = replaced["filter_conditions"]
# SET values and the WHERE clause carry parameter markers; bind them in
# statement order (SET first) and require every parameter to be used
bound, _ = SQLHelper.bind_filter({"u": update_fields, "q": filter_conditions}, parameters)
update_fields, filter_conditions = bound["u"], bound["q"]

# MongoDB update command format
# https://www.mongodb.com/docs/manual/reference/command/update/
Expand Down
Loading
Loading