From f24a7eaaa018aa3fabe354390a74b5e71394e93c Mon Sep 17 00:00:00 2001 From: sdairs Date: Fri, 2 Oct 2026 14:23:06 +0100 Subject: [PATCH 1/3] Add durable Streamlit CSV review desk with atomic approval --- .github/workflows/csv-review-desk.yml | 31 ++ applications/csv-review-desk/.env.example | 6 + applications/csv-review-desk/.gitignore | 7 + .../csv-review-desk/.streamlit/config.toml | 13 + applications/csv-review-desk/README.md | 189 ++++++++++ applications/csv-review-desk/alembic.ini | 3 + applications/csv-review-desk/app.py | 180 ++++++++++ .../csv-review-desk/migrations/env.py | 10 + .../csv-review-desk/migrations/script.py.mako | 13 + .../migrations/versions/20261002_initial.py | 78 ++++ applications/csv-review-desk/pyproject.toml | 4 + .../csv-review-desk/requirements-dev.txt | 45 +++ applications/csv-review-desk/requirements.in | 5 + applications/csv-review-desk/requirements.txt | 41 +++ .../csv-review-desk/reviewdesk/__init__.py | 0 .../csv-review-desk/reviewdesk/database.py | 31 ++ .../csv-review-desk/reviewdesk/models.py | 42 +++ .../csv-review-desk/reviewdesk/service.py | 200 +++++++++++ .../csv-review-desk/reviewdesk/validation.py | 82 +++++ .../csv-review-desk/samples/needs-review.csv | 4 + .../csv-review-desk/samples/valid.csv | 3 + .../csv-review-desk/sql/bootstrap.sql | 7 + applications/csv-review-desk/sql/cleanup.sql | 3 + applications/csv-review-desk/sql/grants.sql | 4 + applications/csv-review-desk/tests/browser.py | 167 +++++++++ .../csv-review-desk/tests/cloud_acceptance.py | 339 ++++++++++++++++++ .../csv-review-desk/tests/persistence.py | 43 +++ .../csv-review-desk/tests/test_validation.py | 68 ++++ 28 files changed, 1618 insertions(+) create mode 100644 .github/workflows/csv-review-desk.yml create mode 100644 applications/csv-review-desk/.env.example create mode 100644 applications/csv-review-desk/.gitignore create mode 100644 applications/csv-review-desk/.streamlit/config.toml create mode 100644 applications/csv-review-desk/README.md create mode 100644 applications/csv-review-desk/alembic.ini create mode 100644 applications/csv-review-desk/app.py create mode 100644 applications/csv-review-desk/migrations/env.py create mode 100644 applications/csv-review-desk/migrations/script.py.mako create mode 100644 applications/csv-review-desk/migrations/versions/20261002_initial.py create mode 100644 applications/csv-review-desk/pyproject.toml create mode 100644 applications/csv-review-desk/requirements-dev.txt create mode 100644 applications/csv-review-desk/requirements.in create mode 100644 applications/csv-review-desk/requirements.txt create mode 100644 applications/csv-review-desk/reviewdesk/__init__.py create mode 100644 applications/csv-review-desk/reviewdesk/database.py create mode 100644 applications/csv-review-desk/reviewdesk/models.py create mode 100644 applications/csv-review-desk/reviewdesk/service.py create mode 100644 applications/csv-review-desk/reviewdesk/validation.py create mode 100644 applications/csv-review-desk/samples/needs-review.csv create mode 100644 applications/csv-review-desk/samples/valid.csv create mode 100644 applications/csv-review-desk/sql/bootstrap.sql create mode 100644 applications/csv-review-desk/sql/cleanup.sql create mode 100644 applications/csv-review-desk/sql/grants.sql create mode 100644 applications/csv-review-desk/tests/browser.py create mode 100644 applications/csv-review-desk/tests/cloud_acceptance.py create mode 100644 applications/csv-review-desk/tests/persistence.py create mode 100644 applications/csv-review-desk/tests/test_validation.py diff --git a/.github/workflows/csv-review-desk.yml b/.github/workflows/csv-review-desk.yml new file mode 100644 index 00000000..e3f4e95c --- /dev/null +++ b/.github/workflows/csv-review-desk.yml @@ -0,0 +1,31 @@ +name: CSV review desk +on: + push: + paths: + - 'applications/csv-review-desk/**' + - '.github/workflows/csv-review-desk.yml' + pull_request: + paths: + - 'applications/csv-review-desk/**' + - '.github/workflows/csv-review-desk.yml' +permissions: + contents: read +jobs: + python-check: + runs-on: ubuntu-latest + defaults: + run: + working-directory: applications/csv-review-desk + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-python@v5 + with: + python-version: '3.12' + cache: pip + cache-dependency-path: applications/csv-review-desk/requirements-dev.txt + - run: python -m pip install -r requirements-dev.txt + - run: python -m pip check + - run: python -m unittest discover -s tests -p 'test_*.py' -v + - run: python -m compileall -q app.py reviewdesk migrations tests + - run: ruff check . + - run: ruff format --check . diff --git a/applications/csv-review-desk/.env.example b/applications/csv-review-desk/.env.example new file mode 100644 index 00000000..efc2b9c0 --- /dev/null +++ b/applications/csv-review-desk/.env.example @@ -0,0 +1,6 @@ +PGHOST=YOUR_DIRECT_POSTGRES_HOSTNAME +PGPORT=5432 +PGDATABASE=postgres +PGUSER=csv_app +PGPASSWORD=YOUR_RUNTIME_PASSWORD +PGSSLROOTCERT=/ABSOLUTE/PATH/TO/postgres-ca.pem diff --git a/applications/csv-review-desk/.gitignore b/applications/csv-review-desk/.gitignore new file mode 100644 index 00000000..a6878249 --- /dev/null +++ b/applications/csv-review-desk/.gitignore @@ -0,0 +1,7 @@ +.env +.deployment/ +.venv/ +__pycache__/ +.pytest_cache/ +.ruff_cache/ +*.png diff --git a/applications/csv-review-desk/.streamlit/config.toml b/applications/csv-review-desk/.streamlit/config.toml new file mode 100644 index 00000000..5a7566bd --- /dev/null +++ b/applications/csv-review-desk/.streamlit/config.toml @@ -0,0 +1,13 @@ +[server] +address = "127.0.0.1" +port = 8501 +headless = true +maxUploadSize = 1 +[theme] +primaryColor = "#17634b" +backgroundColor = "#f7f7f2" +secondaryBackgroundColor = "#eeeeE7" +textColor = "#202c25" +font = "sans serif" +[browser] +gatherUsageStats = false diff --git a/applications/csv-review-desk/README.md b/applications/csv-review-desk/README.md new file mode 100644 index 00000000..0cfcc76e --- /dev/null +++ b/applications/csv-review-desk/README.md @@ -0,0 +1,189 @@ +# CSV review desk + +Stage a CSV, correct its invalid rows, and approve the complete batch into a catalogue backed by [ClickHouse Managed Postgres (public beta)](https://clickhouse.com/docs/products/managed-postgres/overview). This small Streamlit workbench uses SQLAlchemy, psycopg, Alembic and pandas. Invalid values and their feedback remain durable, so closing the browser does not lose a review. + +Two views cover the workflow: upload/list saved batches, then review a batch in a fixed-row data editor. Approval inserts new catalogue items; it never updates an existing SKU. This is a trusted local workbench bound to loopback, without application authentication. Anyone who can reach the process has the same operator authority. Public deployment requires an authentication boundary and appropriate HTTPS/proxy configuration; it is outside this example's tested scope. + +## The rules + +- UTF-8 CSV, **no BOM**, exactly `sku,name,price_cents` in that order. CSV quoting and CRLF/LF are supported. Blank records, extra/missing fields and malformed quoting are rejected. +- Upload **1–262,144 bytes**, **1–200 records**. The raw staging limits are SKU 80, name 240 and price 32 characters, no NUL; exceeding these is a rejected upload/save. Values inside those raw limits can be invalid and corrected later. +- Valid SKU: 1–40 ASCII uppercase letters/digits/underscore/hyphen, beginning with letter/digit. Valid name:1–120 characters, no surrounding whitespace or Unicode control characters (`Cc`). Valid price: decimal integer text **0–1,000,000,000 cents**, no signs, fractional values, exponent notation or leading zeros. Prices are converted to integers only after validation. Currency is not inferred; these are catalogue minor-unit integers. +- Every staged row gets a server-generated UUID and fixed ordinal. The editor hides/disables the UUID; the server independently requires the exact original identity set. Rows cannot be added, removed or moved between batches. +- Save and approval lock the parent batch. Every pending write requires the revision the caller reviewed. A stale write is rejected; saving increments the revision. Reads take a short parent shared lock so header/revision and rows agree. +- Approval validates the saved rows and catalogue again, inserts all items, and marks the batch approved in one transaction. A unique SKU constraint catches a conflict that appears after validation; the entire publication rolls back. Repeating successful approval returns the original time/count even with the original revision. +- Invalid approval commits refreshed validation feedback and an increased revision, with zero items published. After a late constraint conflict, refresh/save to update validation against the winner. Approved batches/rows cannot be edited; triggers guard them and a composite foreign key binds catalogue provenance to the same batch and row. + +The UI re-reads after each write. Only the thread-safe SQLAlchemy engine/pool is shared through `st.cache_resource`; ORM Sessions and query results are not cached. If another session changes the revision, the UI discards old editor state, announces the refresh and suppresses Save/Approve for that rerun. Review the refreshed rows and click again. Unsaved local edits are not durable until saved. + +## 1. Install inside Linux + +Verified with Ubuntu 24.04 arm 64/Python 3.12. The exact tested package versions are in `requirements.in` and all runtime transitive dependencies are pinned in `requirements.txt`. Run the following from this application's directory, in a native Linux filesystem: + +```sh +sudo apt-get update +sudo apt-get install -y --no-install-recommends python3-venv postgresql-client git curl jq openssl +python3 -m venv .venv +source .venv/bin/activate +python -m pip install -r requirements.txt +python -m pip check +``` + +Clone the examples repository if needed and enter `applications/csv-review-desk`. No Node/frontend build is required. `samples/valid.csv` and `samples/needs-review.csv` contain synthetic records; the second needs price/name/duplicate-SKU corrections. Uploaded strings, including formula-looking names, remain literal text. This app never executes formulas or fetches uploaded URLs. + +## 2. Create your own Cloud service + +Install [clickhousectl](https://clickhouse.com/blog/clickhousectl-v0-2-0-postgres-clickpipes-more) and authenticate using an Admin API key from your Cloud organization. Use interactive sign-in so credentials are not shell-history arguments: + +```sh +curl -fsSL https://clickhouse.com/cli | sh +export PATH="$HOME/.local/bin:$PATH" +clickhousectl cloud auth login --interactive +clickhousectl cloud auth status +clickhousectl cloud org list +umask 077 +mkdir -p .deployment +``` + +Create private `.deployment/resources.env` with your organization, available region/size and the direct connection fields supplied after creation: + +```dotenv +CH_ORG_ID=YOUR_ORGANIZATION_UUID +CLOUD_REGION=us-east-1 +PG_SIZE=c6gd.large +``` + +The tested fixture used AWS us-east-1, `c6gd.large`, Postgres 18 and no HA. Check current availability and [pricing](https://clickhouse.com/pricing) before creation. Compute/storage can incur charges; stopping Streamlit does not stop database billing. + +```sh +source .deployment/resources.env +clickhousectl cloud postgres create --org-id "$CH_ORG_ID" \ + --name csv-review-desk-example --provider aws --region "$CLOUD_REGION" \ + --size "$PG_SIZE" --pg-version 18 --ha-type none --json \ + > .deployment/postgres-create.json +PG_SERVICE_ID="$(jq -er '.id' .deployment/postgres-create.json)" +``` + +Save `PG_SERVICE_ID` in `resources.env`. The receipt contains the initial password once. Preserve it privately; `get` does not return passwords. If creation is interrupted, inspect the list before retrying to avoid duplicate billable services. Repeat `get`, not `create`, until state is `running`: + +```sh +clickhousectl cloud postgres get "$PG_SERVICE_ID" --org-id "$CH_ORG_ID" --json +clickhousectl cloud postgres certs get "$PG_SERVICE_ID" --org-id "$CH_ORG_ID" \ + --output .deployment/postgres-ca.pem +``` + +Append `PGHOST=YOUR_DIRECT_HOSTNAME`, `PGPORT=5432`, `PGDATABASE=postgres`, `PGADMIN=postgres` and `PGSSLROOTCERT=/ABSOLUTE/PATH/TO/postgres-ca.pem` to `resources.env`, using actual returned values. Use the direct endpoint for schema work. The CA path must exist in the Linux environment running Python. + +## 3. Separate schema and runtime roles + +Generate passwords once, preserving the file on retries: + +```sh +cat > .deployment/passwords.env <= 1), + row_count integer NOT NULL CHECK(row_count BETWEEN 1 AND 200), + created_at timestamptz NOT NULL DEFAULT now(), approved_at timestamptz, + CONSTRAINT batch_state CHECK((status = 'pending' AND approved_at IS NULL) OR (status = 'approved' AND approved_at IS NOT NULL)) +); +CREATE INDEX batches_recent ON batches(created_at DESC, id); +CREATE TABLE staged_rows ( + id uuid PRIMARY KEY, batch_id uuid NOT NULL REFERENCES batches(id), + position integer NOT NULL CHECK(position BETWEEN 0 AND 199), + sku text NOT NULL CHECK(char_length(sku) <= 80), + name text NOT NULL CHECK(char_length(name) <= 240), + price_cents text NOT NULL CHECK(char_length(price_cents) <= 32), + errors jsonb NOT NULL CHECK(jsonb_typeof(errors) = 'array'), + UNIQUE(batch_id, position), UNIQUE(id, batch_id) +); +CREATE TABLE catalogue_items ( + sku varchar(40) PRIMARY KEY CHECK(sku ~ '^[A-Z0-9][A-Z0-9_-]{0,39}$'), + name varchar(120) NOT NULL CHECK(char_length(name) BETWEEN 1 AND 120 AND btrim(name) = name), + price_cents integer NOT NULL CHECK(price_cents BETWEEN 0 AND 1000000000), + source_batch_id uuid NOT NULL REFERENCES batches(id), source_row_id uuid NOT NULL UNIQUE, + FOREIGN KEY(source_row_id, source_batch_id) REFERENCES staged_rows(id, batch_id) +); +CREATE INDEX catalogue_batch ON catalogue_items(source_batch_id); +CREATE FUNCTION guard_stage() RETURNS trigger LANGUAGE plpgsql AS $$ +DECLARE parent_status text; +BEGIN + IF TG_OP = 'UPDATE' AND (NEW.id, NEW.batch_id, NEW.position) IS DISTINCT FROM (OLD.id, OLD.batch_id, OLD.position) THEN + RAISE EXCEPTION 'Staged row identity is immutable' USING ERRCODE = '23514'; + END IF; + SELECT status INTO parent_status FROM batches WHERE id = COALESCE(NEW.batch_id, OLD.batch_id) FOR UPDATE; + IF parent_status IS DISTINCT FROM 'pending' THEN + RAISE EXCEPTION 'Approved rows are immutable' USING ERRCODE = '23514'; + END IF; + IF TG_OP = 'DELETE' THEN RETURN OLD; ELSE RETURN NEW; END IF; +END $$; +CREATE TRIGGER staged_guard BEFORE INSERT OR UPDATE OR DELETE ON staged_rows FOR EACH ROW EXECUTE FUNCTION guard_stage(); +CREATE FUNCTION guard_batch() RETURNS trigger LANGUAGE plpgsql AS $$ +BEGIN + IF (NEW.id, NEW.row_count, NEW.created_at, NEW.filename) IS DISTINCT FROM (OLD.id, OLD.row_count, OLD.created_at, OLD.filename) + OR NEW.revision <= OLD.revision THEN + RAISE EXCEPTION 'Batch identity or revision invalid' USING ERRCODE = '23514'; + END IF; + IF OLD.status = 'approved' THEN + RAISE EXCEPTION 'Approved batch is immutable' USING ERRCODE = '23514'; + END IF; + IF NEW.status = 'approved' AND ((SELECT count(*) FROM catalogue_items WHERE source_batch_id = NEW.id) <> NEW.row_count + OR (SELECT count(*) FROM staged_rows WHERE batch_id = NEW.id) <> NEW.row_count + OR EXISTS(SELECT 1 FROM staged_rows WHERE batch_id = NEW.id AND errors <> '[]'::jsonb)) THEN + RAISE EXCEPTION 'Approval requires a complete valid publication' USING ERRCODE = '23514'; + END IF; + RETURN NEW; +END $$; +CREATE TRIGGER batch_guard BEFORE UPDATE ON batches FOR EACH ROW EXECUTE FUNCTION guard_batch(); +""") + + +def downgrade(): + op.execute(""" +DROP TABLE catalogue_items; +DROP TABLE staged_rows; +DROP TABLE batches; +DROP FUNCTION guard_stage(); +DROP FUNCTION guard_batch(); +""") diff --git a/applications/csv-review-desk/pyproject.toml b/applications/csv-review-desk/pyproject.toml new file mode 100644 index 00000000..7e9b5362 --- /dev/null +++ b/applications/csv-review-desk/pyproject.toml @@ -0,0 +1,4 @@ +[tool.ruff] +line-length = 100 +[tool.pytest.ini_options] +testpaths = ["tests"] diff --git a/applications/csv-review-desk/requirements-dev.txt b/applications/csv-review-desk/requirements-dev.txt new file mode 100644 index 00000000..30fb970d --- /dev/null +++ b/applications/csv-review-desk/requirements-dev.txt @@ -0,0 +1,45 @@ +alembic==1.20.0 +altair==6.3.0 +anyio==4.15.1 +attrs==26.1.0 +certifi==2026.7.22 +charset-normalizer==3.5.2 +click==8.5.0 +greenlet==3.5.6 +h11==0.16.0 +httptools==0.8.0 +idna==3.20 +itsdangerous==2.2.0 +Jinja2==3.1.6 +jsonschema==4.26.0 +jsonschema-specifications==2025.9.1 +Mako==1.4.3 +MarkupSafe==3.0.3 +narwhals==2.26.0 +numpy==2.5.3 +packaging==26.3 +pandas==3.0.6 +pillow==12.3.0 +playwright==1.55.0 +protobuf==7.36.2 +psycopg==3.3.6 +psycopg-binary==3.3.6 +pyarrow==25.0.1 +pydeck==0.9.3 +pyee==13.0.1 +python-dateutil==2.9.0.post0 +python-multipart==0.0.32 +referencing==0.37.0 +requests==2.34.2 +rpds-py==2026.6.3 +ruff==0.15.7 +six==1.17.0 +SQLAlchemy==2.1.2 +starlette==1.7.0 +streamlit==1.64.0 +toml==0.10.2 +typing_extensions==4.16.0 +urllib3==2.8.0 +uvicorn==0.54.0 +watchdog==6.0.0 +websockets==16.1.1 diff --git a/applications/csv-review-desk/requirements.in b/applications/csv-review-desk/requirements.in new file mode 100644 index 00000000..4969a50f --- /dev/null +++ b/applications/csv-review-desk/requirements.in @@ -0,0 +1,5 @@ +streamlit==1.64.0 +SQLAlchemy==2.1.2 +psycopg[binary]==3.3.6 +alembic==1.20.0 +pandas==3.0.6 diff --git a/applications/csv-review-desk/requirements.txt b/applications/csv-review-desk/requirements.txt new file mode 100644 index 00000000..afe3fc31 --- /dev/null +++ b/applications/csv-review-desk/requirements.txt @@ -0,0 +1,41 @@ +alembic==1.20.0 +altair==6.3.0 +anyio==4.15.1 +attrs==26.1.0 +certifi==2026.7.22 +charset-normalizer==3.5.2 +click==8.5.0 +h11==0.16.0 +httptools==0.8.0 +idna==3.20 +itsdangerous==2.2.0 +Jinja2==3.1.6 +jsonschema==4.26.0 +jsonschema-specifications==2025.9.1 +Mako==1.4.3 +MarkupSafe==3.0.3 +narwhals==2.26.0 +numpy==2.5.3 +packaging==26.3 +pandas==3.0.6 +pillow==12.3.0 +protobuf==7.36.2 +psycopg==3.3.6 +psycopg-binary==3.3.6 +pyarrow==25.0.1 +pydeck==0.9.3 +python-dateutil==2.9.0.post0 +python-multipart==0.0.32 +referencing==0.37.0 +requests==2.34.2 +rpds-py==2026.6.3 +six==1.17.0 +SQLAlchemy==2.1.2 +starlette==1.7.0 +streamlit==1.64.0 +toml==0.10.2 +typing_extensions==4.16.0 +urllib3==2.8.0 +uvicorn==0.54.0 +watchdog==6.0.0 +websockets==16.1.1 diff --git a/applications/csv-review-desk/reviewdesk/__init__.py b/applications/csv-review-desk/reviewdesk/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/applications/csv-review-desk/reviewdesk/database.py b/applications/csv-review-desk/reviewdesk/database.py new file mode 100644 index 00000000..0460fe8c --- /dev/null +++ b/applications/csv-review-desk/reviewdesk/database.py @@ -0,0 +1,31 @@ +import os +from pathlib import Path + +from sqlalchemy import URL, create_engine + + +def make_engine(): + ca = Path(os.environ["PGSSLROOTCERT"]) + if not ca.is_file(): + raise ValueError("PGSSLROOTCERT must point to the downloaded Cloud CA") + url = URL.create( + "postgresql+psycopg", + username=os.environ["PGUSER"], + password=os.environ["PGPASSWORD"], + host=os.environ["PGHOST"], + port=int(os.environ.get("PGPORT", "5432")), + database=os.environ.get("PGDATABASE", "postgres"), + ) + return create_engine( + url, + pool_size=4, + max_overflow=0, + pool_pre_ping=True, + connect_args={ + "sslmode": "verify-full", + "sslrootcert": str(ca.resolve()), + "connect_timeout": 10, + "options": "-c search_path=csv_review,public -c statement_timeout=15000 -c lock_timeout=10000", + }, + hide_parameters=True, + ) diff --git a/applications/csv-review-desk/reviewdesk/models.py b/applications/csv-review-desk/reviewdesk/models.py new file mode 100644 index 00000000..f4af85c4 --- /dev/null +++ b/applications/csv-review-desk/reviewdesk/models.py @@ -0,0 +1,42 @@ +import uuid +from datetime import datetime + +from sqlalchemy import DateTime, ForeignKey, Integer, String, Text, UniqueConstraint, func +from sqlalchemy.dialects.postgresql import JSONB, UUID +from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column + + +class Base(DeclarativeBase): + pass + + +class Batch(Base): + __tablename__ = "batches" + id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), primary_key=True, default=uuid.uuid4) + filename: Mapped[str] = mapped_column(String(120)) + status: Mapped[str] = mapped_column(String(16), default="pending") + revision: Mapped[int] = mapped_column(Integer, default=1) + row_count: Mapped[int] = mapped_column(Integer) + created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now()) + approved_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True)) + + +class StagedRow(Base): + __tablename__ = "staged_rows" + __table_args__ = (UniqueConstraint("batch_id", "position"),) + id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), primary_key=True, default=uuid.uuid4) + batch_id: Mapped[uuid.UUID] = mapped_column(ForeignKey("batches.id")) + position: Mapped[int] = mapped_column(Integer) + sku: Mapped[str] = mapped_column(Text) + name: Mapped[str] = mapped_column(Text) + price_cents: Mapped[str] = mapped_column(Text) + errors: Mapped[list[str]] = mapped_column(JSONB) + + +class CatalogueItem(Base): + __tablename__ = "catalogue_items" + sku: Mapped[str] = mapped_column(String(40), primary_key=True) + name: Mapped[str] = mapped_column(String(120)) + price_cents: Mapped[int] = mapped_column(Integer) + source_batch_id: Mapped[uuid.UUID] = mapped_column(ForeignKey("batches.id")) + source_row_id: Mapped[uuid.UUID] = mapped_column(ForeignKey("staged_rows.id"), unique=True) diff --git a/applications/csv-review-desk/reviewdesk/service.py b/applications/csv-review-desk/reviewdesk/service.py new file mode 100644 index 00000000..39dece12 --- /dev/null +++ b/applications/csv-review-desk/reviewdesk/service.py @@ -0,0 +1,200 @@ +import uuid +from datetime import datetime, timezone + +from sqlalchemy import func, select +from sqlalchemy.exc import IntegrityError +from sqlalchemy.orm import Session + +from reviewdesk.models import Batch, CatalogueItem, StagedRow +from reviewdesk.validation import HEADERS, ReviewError, parse_csv, raw_record, validate_records + + +def _uuid(value): + try: + return uuid.UUID(str(value)) + except (ValueError, TypeError, AttributeError) as exc: + raise ReviewError("Unknown batch or row.") from exc + + +def _locked(session, batch_id, read=False): + batch = session.scalar( + select(Batch).where(Batch.id == _uuid(batch_id)).with_for_update(read=read) + ) + if batch is None: + raise ReviewError("Batch not found.") + return batch + + +def _revision(batch, expected): + if type(expected) is not int or batch.revision != expected: + raise ReviewError( + "This batch changed in another session. Refresh and review the committed rows before saving." + ) + if batch.status != "pending": + raise ReviewError("Approved batches cannot be edited.") + + +def _rows(session, batch): + return list( + session.scalars( + select(StagedRow).where(StagedRow.batch_id == batch.id).order_by(StagedRow.position) + ) + ) + + +def _feedback(session, rows): + records = [{key: getattr(row, key) for key in HEADERS} for row in rows] + existing = session.scalars( + select(CatalogueItem.sku).where(CatalogueItem.sku.in_([r["sku"] for r in records])) + ) + errors = validate_records(records, existing) + for row, issues in zip(rows, errors): + row.errors = issues + return errors + + +def create_batch(engine, filename, data): + records = parse_csv(data) + filename = str(filename).replace("\x00", "")[:120] or "upload.csv" + with Session(engine) as session, session.begin(): + batch = Batch(filename=filename, row_count=len(records)) + session.add(batch) + session.flush() + rows = [ + StagedRow(batch_id=batch.id, position=i, errors=[], **record) + for i, record in enumerate(records) + ] + session.add_all(rows) + _feedback(session, rows) + session.flush() + result = str(batch.id) + return result + + +def save_rows(engine, batch_id, expected_revision, submitted): + with Session(engine) as session, session.begin(): + batch = _locked(session, batch_id) + _revision(batch, expected_revision) + rows = _rows(session, batch) + if not isinstance(submitted, list) or len(submitted) != len(rows): + raise ReviewError( + "Save the complete fixed set of rows; rows cannot be added or removed." + ) + values = {} + for record in submitted: + if not isinstance(record, dict) or set(record) != {"id", *HEADERS}: + raise ReviewError("Submit the fixed row identity and its three editable fields.") + row_id = _uuid(record["id"]) + if row_id in values: + raise ReviewError("Repeated row identity.") + values[row_id] = raw_record({key: record[key] for key in HEADERS}) + if set(values) != {row.id for row in rows}: + raise ReviewError("The submitted row identities do not belong to this batch.") + for row in rows: + for key, value in values[row.id].items(): + setattr(row, key, value) + _feedback(session, rows) + batch.revision += 1 + return get_batch(engine, batch_id) + + +def approve_batch(engine, batch_id, expected_revision): + failure = None + try: + with Session(engine) as session, session.begin(): + batch = _locked(session, batch_id) + if batch.status == "approved": + return { + "batch_id": str(batch.id), + "published_rows": batch.row_count, + "approved_at": batch.approved_at.isoformat(), + } + _revision(batch, expected_revision) + rows = _rows(session, batch) + errors = _feedback(session, rows) + if any(errors): + batch.revision += 1 + failure = ( + "Approval stopped: correct the saved row issues, then review and approve again." + ) + else: + for row in sorted(rows, key=lambda row: row.sku): + session.add( + CatalogueItem( + sku=row.sku, + name=row.name, + price_cents=int(row.price_cents), + source_batch_id=batch.id, + source_row_id=row.id, + ) + ) + session.flush() + batch.status = "approved" + batch.approved_at = datetime.now(timezone.utc) + batch.revision += 1 + result = { + "batch_id": str(batch.id), + "published_rows": batch.row_count, + "approved_at": batch.approved_at.isoformat(), + } + except IntegrityError as exc: + # The transaction rolls back every inserted row if another batch wins a SKU. + if ( + getattr(getattr(exc.orig, "diag", None), "constraint_name", None) + == "catalogue_items_pkey" + ): + raise ReviewError( + "Approval rolled back: a catalogue SKU conflict appeared. Refresh and save to update validation." + ) from exc + raise ReviewError( + "Approval rolled back because a database constraint rejected the publication." + ) from exc + if failure: + raise ReviewError(failure) + return result + + +def get_batch(engine, batch_id): + with Session(engine) as session, session.begin(): + # Keep the parent revision and its rows in one consistent shared-lock read. + batch = _locked(session, batch_id, read=True) + rows = _rows(session, batch) + return { + "id": str(batch.id), + "filename": batch.filename, + "status": batch.status, + "revision": batch.revision, + "row_count": batch.row_count, + "approved_at": batch.approved_at.isoformat() if batch.approved_at else None, + "rows": [ + { + "id": str(row.id), + **{key: getattr(row, key) for key in HEADERS}, + "errors": list(row.errors), + } + for row in rows + ], + } + + +def list_batches(engine, page=1): + if type(page) is not int or page < 1: + raise ReviewError("Page must be a positive integer.") + with Session(engine) as session: + total = session.scalar(select(func.count()).select_from(Batch)) + batches = session.scalars( + select(Batch) + .order_by(Batch.created_at.desc(), Batch.id) + .limit(20) + .offset((page - 1) * 20) + ) + return total, [ + { + "id": str(b.id), + "filename": b.filename, + "status": b.status, + "revision": b.revision, + "row_count": b.row_count, + } + for b in batches + ] diff --git a/applications/csv-review-desk/reviewdesk/validation.py b/applications/csv-review-desk/reviewdesk/validation.py new file mode 100644 index 00000000..15f984cc --- /dev/null +++ b/applications/csv-review-desk/reviewdesk/validation.py @@ -0,0 +1,82 @@ +import csv +import io +import re +import unicodedata +from collections import Counter + +MAX_BYTES = 262_144 +MAX_ROWS = 200 +HEADERS = ["sku", "name", "price_cents"] +RAW_LIMITS = {"sku": 80, "name": 240, "price_cents": 32} +SKU = re.compile(r"[A-Z0-9][A-Z0-9_-]{0,39}", re.ASCII) +PRICE = re.compile(r"0|[1-9][0-9]{0,9}", re.ASCII) + + +class ReviewError(Exception): + pass + + +def raw_record(record): + if set(record) != set(HEADERS): + raise ReviewError("Each row must contain sku, name and price_cents.") + result = {} + for field, limit in RAW_LIMITS.items(): + value = record[field] + if not isinstance(value, str) or len(value) > limit or "\x00" in value: + raise ReviewError(f"{field} must be text of at most {limit} characters, without NUL.") + result[field] = value + return result + + +def parse_csv(data): + if not isinstance(data, bytes) or not data or len(data) > MAX_BYTES: + raise ReviewError("Upload 1–262,144 bytes.") + try: + text = data.decode("utf-8") + reader = csv.reader(io.StringIO(text, newline=""), strict=True) + if next(reader, None) != HEADERS: + raise ReviewError("Use exactly the UTF-8 header sku,name,price_cents (no BOM).") + rows = [] + for values in reader: + if len(values) != 3: + raise ReviewError( + "Every CSV record must have exactly three fields; blank records are rejected." + ) + rows.append(raw_record(dict(zip(HEADERS, values)))) + if len(rows) > MAX_ROWS: + raise ReviewError("A batch can contain at most 200 records.") + if not rows: + raise ReviewError("Include at least one record below the header.") + return rows + except (UnicodeDecodeError, csv.Error) as exc: + raise ReviewError("Use well-formed UTF-8 CSV with valid quoting.") from exc + + +def validate_records(records, existing=()): + counts = Counter(r["sku"] for r in records) + existing = set(existing) + errors = [] + for row in records: + issues = [] + if not SKU.fullmatch(row["sku"]): + issues.append( + "SKU: 1–40 uppercase letters, digits, underscore or hyphen; begin with a letter/digit." + ) + if ( + not 1 <= len(row["name"]) <= 120 + or row["name"].strip() != row["name"] + or any(unicodedata.category(c) == "Cc" for c in row["name"]) + ): + issues.append( + "Name: 1–120 characters, no surrounding whitespace or control characters." + ) + if not PRICE.fullmatch(row["price_cents"]) or int(row["price_cents"]) > 1_000_000_000: + issues.append( + "Price: whole cents from 0 to 1,000,000,000; no signs, decimals or leading zeros." + ) + if counts[row["sku"]] > 1: + issues.append("SKU is repeated in this batch.") + if row["sku"] in existing: + issues.append("SKU already exists in the catalogue; choose a different SKU.") + errors.append(issues) + return errors diff --git a/applications/csv-review-desk/samples/needs-review.csv b/applications/csv-review-desk/samples/needs-review.csv new file mode 100644 index 00000000..a3821714 --- /dev/null +++ b/applications/csv-review-desk/samples/needs-review.csv @@ -0,0 +1,4 @@ +sku,name,price_cents +NOTE-01,Notebook,3.50 +MUG-02,,1299 +NOTE-01,Notebook duplicate,350 diff --git a/applications/csv-review-desk/samples/valid.csv b/applications/csv-review-desk/samples/valid.csv new file mode 100644 index 00000000..d73c0a88 --- /dev/null +++ b/applications/csv-review-desk/samples/valid.csv @@ -0,0 +1,3 @@ +sku,name,price_cents +DESK-01,Desk tray,1299 +PEN-02,Pen set,450 diff --git a/applications/csv-review-desk/sql/bootstrap.sql b/applications/csv-review-desk/sql/bootstrap.sql new file mode 100644 index 00000000..0f4b52fd --- /dev/null +++ b/applications/csv-review-desk/sql/bootstrap.sql @@ -0,0 +1,7 @@ +\getenv migrator_password PG_MIGRATION_PASSWORD +\getenv app_password PG_APP_PASSWORD +CREATE ROLE csv_migrator LOGIN PASSWORD :'migrator_password' NOSUPERUSER NOCREATEDB NOCREATEROLE; +CREATE ROLE csv_app LOGIN PASSWORD :'app_password' NOSUPERUSER NOCREATEDB NOCREATEROLE; +CREATE SCHEMA csv_review AUTHORIZATION csv_migrator; +REVOKE CREATE ON SCHEMA public FROM PUBLIC; +SELECT format('REVOKE CREATE ON DATABASE %I FROM PUBLIC', current_database()) \gexec diff --git a/applications/csv-review-desk/sql/cleanup.sql b/applications/csv-review-desk/sql/cleanup.sql new file mode 100644 index 00000000..3f8ee0d6 --- /dev/null +++ b/applications/csv-review-desk/sql/cleanup.sql @@ -0,0 +1,3 @@ +DROP SCHEMA csv_review CASCADE; +DROP ROLE csv_app; +DROP ROLE csv_migrator; diff --git a/applications/csv-review-desk/sql/grants.sql b/applications/csv-review-desk/sql/grants.sql new file mode 100644 index 00000000..93ede203 --- /dev/null +++ b/applications/csv-review-desk/sql/grants.sql @@ -0,0 +1,4 @@ +GRANT USAGE ON SCHEMA csv_review TO csv_app; +GRANT SELECT, INSERT, UPDATE ON csv_review.batches, csv_review.staged_rows TO csv_app; +GRANT SELECT, INSERT ON csv_review.catalogue_items TO csv_app; +REVOKE ALL ON csv_review.alembic_version FROM csv_app; diff --git a/applications/csv-review-desk/tests/browser.py b/applications/csv-review-desk/tests/browser.py new file mode 100644 index 00000000..43888f7e --- /dev/null +++ b/applications/csv-review-desk/tests/browser.py @@ -0,0 +1,167 @@ +"""Real Streamlit upload/editor workflow; run on the dedicated local workbench.""" + +import os +import time +from pathlib import Path +from urllib.parse import parse_qs, urlparse + +from playwright.sync_api import expect, sync_playwright +from sqlalchemy import text + +from reviewdesk.database import make_engine +from reviewdesk.service import get_batch + +expect.set_options(timeout=30000) +BASE_URL = "http://127.0.0.1:8501" +ARTIFACTS = Path(os.environ.get("EVIDENCE_DIR", "/tmp/csv-review-desk-evidence")) +ARTIFACTS.mkdir(parents=True, exist_ok=True) + + +def ready(page): + expect(page.get_by_role("button", name="Save corrections", exact=True)).to_be_visible() + expect(page.locator('[data-testid="data-grid-canvas"]').first).to_be_visible() + + +def edit_cell(page, column, value): + canvas = page.locator('[data-testid="data-grid-canvas"]').first + box = canvas.bounding_box() + page.mouse.click(box["x"] + 25, box["y"] + 54) + for _ in range(column): + page.keyboard.press("ArrowRight") + page.keyboard.press("Enter") + editor = page.get_by_role("textbox").last + expect(editor).to_be_visible() + editor.fill(value) + editor.press("Tab") + # Streamlit debounces editor widget-state serialization. Let that bounded + # frontend commit settle before submitting the enclosing form. + page.wait_for_timeout(750) + + +def run(): + engine = make_engine() + errors = [] + uploads = [] + with sync_playwright() as playwright: + browser = playwright.chromium.launch(headless=True) + first_context = browser.new_context(viewport={"width": 1440, "height": 1100}) + second_context = browser.new_context(viewport={"width": 1440, "height": 1100}) + first, second = first_context.new_page(), second_context.new_page() + first_context.tracing.start(screenshots=True, snapshots=True) + for page in [first, second]: + page.on("pageerror", lambda error: errors.append(str(error))) + first.on( + "request", + lambda request: uploads.append(request.url) if "upload_file" in request.url else None, + ) + try: + first.goto(BASE_URL) + expect(first.get_by_role("heading", name="CSV review desk", exact=True)).to_be_visible() + first.locator('input[type="file"]').set_input_files( + { + "name": "notebook-review.csv", + "mimeType": "text/csv", + "buffer": f"sku,name,price_cents\nBROWSER-{int(time.time())},Notebook,3.50\n".encode(), + } + ) + first.get_by_role("button", name="Upload for review", exact=True).click() + expect(first.get_by_role("heading", name="Review batch", exact=True)).to_be_visible() + ready(first) + expect(first.get_by_role("button", name="Approve batch", exact=True)).to_be_disabled() + batch_id = parse_qs(urlparse(first.url).query)["batch"][0] + first.screenshot(path=str(ARTIFACTS / "csv-invalid.png"), full_page=True) + edit_cell(first, 2, "350") + first.get_by_role("button", name="Save corrections", exact=True).click() + expect( + first.get_by_text( + "Corrections saved. Review the committed validation below.", exact=True + ) + ).to_be_visible() + ready(first) + self_check = get_batch(engine, batch_id) + assert self_check["revision"] == 2 + assert self_check["rows"][0]["price_cents"] == "350" + assert self_check["rows"][0]["errors"] == [] + expect(first.get_by_role("button", name="Approve batch", exact=True)).to_be_enabled() + + # Another actual browser session changes the rows after first has reviewed revision2. + second.goto(first.url) + ready(second) + edit_cell(second, 1, "Notebook revised") + second.get_by_role("button", name="Save corrections", exact=True).click() + expect( + second.get_by_text( + "Corrections saved. Review the committed validation below.", exact=True + ) + ).to_be_visible() + ready(second) + assert get_batch(engine, batch_id)["revision"] == 3 + first.get_by_role("button", name="Approve batch", exact=True).click() + expect( + first.get_by_text( + "The saved revision changed. The editor now shows committed rows; unsaved local changes were discarded.", + exact=True, + ) + ).to_be_visible() + assert get_batch(engine, batch_id)["status"] == "pending" + with engine.connect() as conn: + assert ( + conn.scalar( + text("SELECT count(*) FROM catalogue_items WHERE source_batch_id=:id"), + {"id": batch_id}, + ) + == 0 + ) + expect(first.locator('[data-testid="glide-cell-1-0"]')).to_have_text("Notebook revised") + first.screenshot(path=str(ARTIFACTS / "csv-desktop.png"), full_page=True) + + # Fresh click after inspecting refreshed revision3 can publish it. + first.get_by_role("button", name="Approve batch", exact=True).click() + expect(first.get_by_text("Approved: 1 catalogue items.", exact=True)).to_be_visible() + current = get_batch(engine, batch_id) + assert current["status"] == "approved" and current["revision"] == 4 + assert current["rows"][0]["name"] == "Notebook revised" + approved_at = current["approved_at"] + first.get_by_role("button", name="Confirm approval again", exact=True).click() + expect( + first.get_by_text( + "Already approved: 1 items; original approval retained.", exact=True + ) + ).to_be_visible() + assert get_batch(engine, batch_id)["approved_at"] == approved_at + first.screenshot(path=str(ARTIFACTS / "csv-approved.png"), full_page=True) + + # Native upload endpoint refuses a request without an XSRF header. + assert uploads, "Actual upload endpoint was observed" + response = first_context.request.put( + uploads[0], + multipart={ + "file": { + "name": "no-token.csv", + "mimeType": "text/csv", + "buffer": b"sku,name,price_cents\nA,a,1\n", + } + }, + ) + assert response.status == 403, f"Expected native XSRF403, got {response.status}" + mobile = browser.new_page(viewport={"width": 390, "height": 844}) + mobile.goto(first.url) + expect( + mobile.get_by_role("button", name="Confirm approval again", exact=True) + ).to_be_visible() + assert mobile.evaluate("document.documentElement.scrollWidth <= window.innerWidth") + mobile.screenshot(path=str(ARTIFACTS / "csv-mobile.png"), full_page=True) + assert not errors, errors + (ARTIFACTS / "persistence.json").write_text(__import__("json").dumps(current)) + print( + "PASS actual upload, invalid feedback, editor correction, durable exact price, two-session stale-click rejection, fresh approval, repeat result, native XSRF403, desktop/mobile fit, no page errors" + ) + print("Batch", batch_id) + finally: + first_context.tracing.stop(path=str(ARTIFACTS / "browser-trace.zip")) + browser.close() + engine.dispose() + + +if __name__ == "__main__": + run() diff --git a/applications/csv-review-desk/tests/cloud_acceptance.py b/applications/csv-review-desk/tests/cloud_acceptance.py new file mode 100644 index 00000000..80990298 --- /dev/null +++ b/applications/csv-review-desk/tests/cloud_acceptance.py @@ -0,0 +1,339 @@ +"""Run only against this example's dedicated Cloud database; migrator cleans owned fixtures.""" + +import os +import socket +import subprocess +import tempfile +import threading +import time +import unittest +import uuid +from concurrent.futures import ThreadPoolExecutor + +import psycopg +from sqlalchemy import event, text +from sqlalchemy.exc import DBAPIError +from unittest.mock import patch + +from reviewdesk.database import make_engine +from reviewdesk.service import approve_batch, create_batch, get_batch, save_rows +from reviewdesk.validation import HEADERS, ReviewError + + +class CloudAcceptance(unittest.TestCase): + @classmethod + def setUpClass(cls): + cls.engine = make_engine() + cls.batch_ids = [] + cls.prefix = "T" + uuid.uuid4().hex[:10].upper() + with cls.engine.connect() as conn: + print("Postgres", conn.scalar(text("SELECT version()"))) + print("Runtime", conn.scalar(text("SELECT current_user"))) + + @classmethod + def tearDownClass(cls): + # This cleanup authority is never passed to the Streamlit server. + if cls.batch_ids: + with psycopg.connect( + user="csv_migrator", + password=os.environ["TEST_MIGRATOR_PASSWORD"], + options="-c search_path=csv_review,public", + sslmode="verify-full", + ) as conn: + conn.execute("ALTER TABLE staged_rows DISABLE TRIGGER staged_guard") + conn.execute( + "DELETE FROM catalogue_items WHERE source_batch_id = ANY(%s)", (cls.batch_ids,) + ) + conn.execute("DELETE FROM staged_rows WHERE batch_id = ANY(%s)", (cls.batch_ids,)) + conn.execute("DELETE FROM batches WHERE id = ANY(%s)", (cls.batch_ids,)) + conn.execute("ALTER TABLE staged_rows ENABLE TRIGGER staged_guard") + cls.engine.dispose() + + def batch(self, suffix, rows): + csv = ( + "sku,name,price_cents\n" + + "\n".join(f"{self.prefix}-{suffix}-{sku},{name},{price}" for sku, name, price in rows) + + "\n" + ) + batch_id = create_batch(self.engine, suffix + ".csv", csv.encode()) + self.batch_ids.append(uuid.UUID(batch_id)) + return get_batch(self.engine, batch_id) + + def payload(self, batch): + return [{key: row[key] for key in ["id", *HEADERS]} for row in batch["rows"]] + + def count(self, batch_id): + with self.engine.connect() as conn: + return conn.scalar( + text("SELECT count(*) FROM catalogue_items WHERE source_batch_id=:id"), + {"id": uuid.UUID(batch_id)}, + ) + + def outcome(self, action): + try: + return ("ok", action()) + except ReviewError as exc: + return ("rejected", str(exc)) + + def wait_for_two_locks(self): + deadline = time.monotonic() + 7 + while time.monotonic() < deadline: + with self.engine.connect() as conn: + count = conn.scalar( + text( + "SELECT count(*) FROM pg_stat_activity WHERE usename=current_user AND wait_event_type='Lock' AND query ILIKE '%batches%'" + ) + ) + if count >= 2: + return + time.sleep(0.1) + self.fail("Both independent transactions must be observed waiting for the held parent row") + + def test_invalid_rows_correction_stable_identity_exact_price(self): + batch = self.batch("CORRECT", [("A", "", "3.50"), ("A", "Notebook", "350")]) + ids = [r["id"] for r in batch["rows"]] + self.assertTrue(all(r["errors"] for r in batch["rows"])) + rows = self.payload(batch) + rows[0].update(name="Notebook", price_cents="350") + rows[1]["sku"] += "-SECOND" + corrected = save_rows(self.engine, batch["id"], 1, list(reversed(rows))) + self.assertEqual([r["id"] for r in corrected["rows"]], ids) + self.assertFalse(any(r["errors"] for r in corrected["rows"])) + result = approve_batch(self.engine, batch["id"], 2) + self.assertEqual(result["published_rows"], 2) + with self.engine.connect() as conn: + self.assertEqual( + conn.scalar( + text("SELECT sum(price_cents) FROM catalogue_items WHERE source_batch_id=:id"), + {"id": uuid.UUID(batch["id"])}, + ), + 700, + ) + + def test_stale_save_and_forged_row_identity_rejected(self): + batch = self.batch("STALE", [("A", "Original", "10")]) + payload = self.payload(batch) + payload[0]["name"] = "Committed" + save_rows(self.engine, batch["id"], 1, payload) + payload[0]["name"] = "Stale overwrite" + with self.assertRaisesRegex(ReviewError, "changed"): + save_rows(self.engine, batch["id"], 1, payload) + with self.assertRaisesRegex(ReviewError, "changed"): + approve_batch(self.engine, batch["id"], 1) + payload[0]["id"] = str(uuid.uuid4()) + with self.assertRaisesRegex(ReviewError, "identities"): + save_rows(self.engine, batch["id"], 2, payload) + self.assertEqual(get_batch(self.engine, batch["id"])["rows"][0]["name"], "Committed") + self.assertEqual(self.count(batch["id"]), 0) + + def test_simultaneous_approvals_return_original_result(self): + batch = self.batch("RACE", [("A", "One", "10"), ("B", "Two", "20")]) + barrier = threading.Barrier(3) + + def action(): + barrier.wait() + return approve_batch(self.engine, batch["id"], 1) + + with ThreadPoolExecutor(max_workers=2) as pool: + with self.engine.begin() as holder: + holder.execute( + text("SELECT id FROM batches WHERE id=:id FOR UPDATE"), + {"id": uuid.UUID(batch["id"])}, + ) + futures = [pool.submit(action) for _ in range(2)] + barrier.wait() + self.wait_for_two_locks() + results = [future.result(timeout=25) for future in futures] + self.assertEqual(results[0], results[1]) + self.assertEqual(self.count(batch["id"]), 2) + self.assertEqual(approve_batch(self.engine, batch["id"], 1), results[0]) + + def test_edit_versus_approve_ordering(self): + batch = self.batch("ORDER", [("A", "Before", "10")]) + payload = self.payload(batch) + payload[0].update(name="After", price_cents="99") + barrier = threading.Barrier(3) + + def save(): + barrier.wait() + return self.outcome(lambda: save_rows(self.engine, batch["id"], 1, payload)) + + def approve(): + barrier.wait() + return self.outcome(lambda: approve_batch(self.engine, batch["id"], 1)) + + with ThreadPoolExecutor(max_workers=2) as pool: + with self.engine.begin() as holder: + holder.execute( + text("SELECT id FROM batches WHERE id=:id FOR UPDATE"), + {"id": uuid.UUID(batch["id"])}, + ) + edit_future, approve_future = pool.submit(save), pool.submit(approve) + barrier.wait() + self.wait_for_two_locks() + edit, approval = edit_future.result(timeout=25), approve_future.result(timeout=25) + self.assertEqual(sorted([edit[0], approval[0]]), ["ok", "rejected"]) + current = get_batch(self.engine, batch["id"]) + if approval[0] == "ok": + self.assertEqual(current["status"], "approved") + self.assertEqual(current["rows"][0]["name"], "Before") + self.assertEqual(self.count(batch["id"]), 1) + else: + self.assertEqual(current["status"], "pending") + self.assertEqual(current["rows"][0]["name"], "After") + self.assertEqual(self.count(batch["id"]), 0) + + def test_catalogue_conflict_after_validation_rolls_back_all(self): + batch = self.batch("LATE", [("A", "New", "10"), ("Z", "Shared", "20")]) + other = self.batch("WINNER", [("Z", "Winner", "99")]) + rows = self.payload(other) + rows[0]["sku"] = batch["rows"][1]["sku"] + save_rows(self.engine, other["id"], 1, rows) + paused, release = threading.Event(), threading.Event() + competitor_engine = make_engine() + + def pause_insert(conn, cursor, statement, parameters, context, executemany): + if statement.startswith("INSERT INTO catalogue_items"): + paused.set() + if not release.wait(30): + raise AssertionError("Competing publication did not release insert") + + event.listen(self.engine, "before_cursor_execute", pause_insert) + try: + with ThreadPoolExecutor(max_workers=1) as pool: + future = pool.submit( + self.outcome, lambda: approve_batch(self.engine, batch["id"], 1) + ) + self.assertTrue( + paused.wait(30), "Approval reached insert after its validation read" + ) + approve_batch(competitor_engine, other["id"], 2) + release.set() + outcome = future.result(timeout=25) + self.assertEqual(outcome[0], "rejected") + self.assertIn("SKU conflict", outcome[1]) + finally: + release.set() + event.remove(self.engine, "before_cursor_execute", pause_insert) + competitor_engine.dispose() + self.assertEqual(self.count(batch["id"]), 0) + self.assertEqual(get_batch(self.engine, batch["id"])["status"], "pending") + with self.engine.connect() as conn: + self.assertEqual( + conn.scalar( + text("SELECT name FROM catalogue_items WHERE sku=:sku"), {"sku": rows[0]["sku"]} + ), + "Winner", + ) + + def test_invalid_approval_feedback_persists_without_partial_publish(self): + batch = self.batch("INVALID", [("A", "Good", "10"), ("B", "", "1.0")]) + with self.assertRaisesRegex(ReviewError, "correct"): + approve_batch(self.engine, batch["id"], 1) + current = get_batch(self.engine, batch["id"]) + self.assertEqual(current["revision"], 2) + self.assertTrue(current["rows"][1]["errors"]) + self.assertEqual(self.count(batch["id"]), 0) + + def test_existing_catalogue_validation_and_immutable_approved_rows(self): + batch = self.batch("EXIST", [("A", "First", "10")]) + approve_batch(self.engine, batch["id"], 1) + other = self.batch("EXIST2", [("A", "Second", "20")]) + payload = self.payload(other) + payload[0]["sku"] = batch["rows"][0]["sku"] + corrected = save_rows(self.engine, other["id"], 1, payload) + self.assertIn("already exists", corrected["rows"][0]["errors"][0]) + with self.assertRaises(ReviewError): + save_rows(self.engine, batch["id"], 2, self.payload(batch)) + with self.assertRaises(DBAPIError) as failure, self.engine.begin() as conn: + conn.execute( + text("UPDATE staged_rows SET name='Changed' WHERE batch_id=:id"), + {"id": uuid.UUID(batch["id"])}, + ) + self.assertEqual(failure.exception.orig.sqlstate, "23514") + + def test_database_constraints_and_restricted_runtime(self): + batch = self.batch("DB", [("A", "One", "10")]) + other = self.batch("DB2", [("A", "Other", "10")]) + with self.assertRaises(DBAPIError) as failure, self.engine.begin() as conn: + conn.execute( + text("INSERT INTO catalogue_items VALUES ('MISMATCH','Wrong',1,:batch,:row)"), + {"batch": uuid.UUID(other["id"]), "row": uuid.UUID(batch["rows"][0]["id"])}, + ) + self.assertEqual(failure.exception.orig.sqlstate, "23503") + with self.assertRaises(DBAPIError) as failure, self.engine.begin() as conn: + conn.execute( + text( + "UPDATE batches SET status='approved', approved_at=now(), revision=revision+1 WHERE id=:id" + ), + {"id": uuid.UUID(batch["id"])}, + ) + self.assertEqual(failure.exception.orig.sqlstate, "23514") + for sql in [ + "CREATE TABLE csv_review.forbidden(id integer)", + "CREATE TABLE public.forbidden(id integer)", + "SELECT * FROM csv_review.alembic_version", + "DELETE FROM catalogue_items", + "UPDATE catalogue_items SET name='Changed'", + ]: + with ( + self.subTest(sql=sql), + self.assertRaises(DBAPIError) as failure, + self.engine.begin() as conn, + ): + conn.execute(text(sql)) + self.assertEqual(failure.exception.orig.sqlstate, "42501") + + def test_driver_tls_verification(self): + with self.engine.connect() as conn: + self.assertTrue( + conn.scalar(text("SELECT ssl FROM pg_stat_ssl WHERE pid=pg_backend_pid()")) + ) + with tempfile.TemporaryDirectory() as directory: + ca = directory + "/wrong.pem" + subprocess.run( + [ + "openssl", + "req", + "-x509", + "-newkey", + "rsa:2048", + "-nodes", + "-days", + "1", + "-subj", + "/CN=Wrong CA", + "-keyout", + directory + "/key", + "-out", + ca, + ], + check=True, + stdout=subprocess.DEVNULL, + stderr=subprocess.DEVNULL, + ) + with patch.dict(os.environ, {"PGSSLROOTCERT": ca}): + wrong = make_engine() + try: + with self.assertRaises(DBAPIError) as failure: + wrong.connect() + self.assertIn("certificate verify failed", str(failure.exception.orig)) + finally: + wrong.dispose() + hostaddr = socket.gethostbyname(os.environ["PGHOST"]) + with self.assertRaises(psycopg.OperationalError) as failure: + psycopg.connect( + host="hostname-mismatch.invalid", + hostaddr=hostaddr, + user=os.environ["PGUSER"], + password=os.environ["PGPASSWORD"], + dbname=os.environ["PGDATABASE"], + sslmode="verify-full", + sslrootcert=os.environ["PGSSLROOTCERT"], + connect_timeout=10, + ) + self.assertRegex(str(failure.exception), "certificate.*does not match host name") + + +if __name__ == "__main__": + unittest.main(verbosity=2) diff --git a/applications/csv-review-desk/tests/persistence.py b/applications/csv-review-desk/tests/persistence.py new file mode 100644 index 00000000..dae7d1cd --- /dev/null +++ b/applications/csv-review-desk/tests/persistence.py @@ -0,0 +1,43 @@ +"""After an actual process restart, compare the browser's durable approved batch.""" + +import json +import os +from pathlib import Path + +from playwright.sync_api import expect, sync_playwright + +from reviewdesk.database import make_engine +from reviewdesk.service import get_batch + + +def run(): + evidence = Path(os.environ.get("EVIDENCE_DIR", "/tmp/csv-review-desk-evidence")) + before = json.loads((evidence / "persistence.json").read_text()) + engine = make_engine() + try: + after = get_batch(engine, before["id"]) + assert before == after, ( + "Saved revision, identities, raw rows and original approval must survive" + ) + with sync_playwright() as playwright: + browser = playwright.chromium.launch() + page = browser.new_page(viewport={"width": 1440, "height": 1100}) + page.goto("http://127.0.0.1:8501/?batch=" + before["id"]) + expect( + page.get_by_role("button", name="Confirm approval again", exact=True) + ).to_be_visible(timeout=30000) + expect(page.locator('[data-testid="data-grid-canvas"]').first).to_be_visible() + expect(page.locator('[data-testid="glide-cell-1-0"]')).to_have_text( + before["rows"][0]["name"] + ) + browser.close() + print( + "PASS real restarted process reads identical approved batch and displays it at its durable URL" + ) + print("Batch", before["id"], "revision", before["revision"], "rows", before["row_count"]) + finally: + engine.dispose() + + +if __name__ == "__main__": + run() diff --git a/applications/csv-review-desk/tests/test_validation.py b/applications/csv-review-desk/tests/test_validation.py new file mode 100644 index 00000000..44dcc0df --- /dev/null +++ b/applications/csv-review-desk/tests/test_validation.py @@ -0,0 +1,68 @@ +import unittest + +from reviewdesk.validation import MAX_BYTES, ReviewError, parse_csv, validate_records + + +class CSVValidationTests(unittest.TestCase): + def test_quoted_utf8_and_exact_integer_text(self): + rows = parse_csv( + 'sku,name,price_cents\r\nA,"Café, tray",1299\r\nB,Free sample,0\r\n'.encode() + ) + self.assertEqual(rows[0]["name"], "Café, tray") + self.assertEqual(rows[0]["price_cents"], "1299") + self.assertEqual(validate_records(rows), [[], []]) + + def test_exact_header_and_encoding(self): + for value in [ + b"sku,name,price\nA,a,1\n", + b"name,sku,price_cents\n", + b"\xef\xbb\xbfsku,name,price_cents\nA,a,1\n", + b"sku,name,price_cents\nA,\xff,1\n", + ]: + with self.subTest(value=value), self.assertRaises(ReviewError): + parse_csv(value) + + def test_shape_quoting_and_bounds(self): + for value in [ + b"", + b"sku,name,price_cents\n", + b"sku,name,price_cents\nA,a\n", + b"sku,name,price_cents\n\n", + b'sku,name,price_cents\nA,"broken,1\n', + b"x" * (MAX_BYTES + 1), + b"sku,name,price_cents\n" + b"A,a,1\n" * 201, + b"sku,name,price_cents\nA," + b"a" * 241 + b",1\n", + ]: + with self.subTest(length=len(value)), self.assertRaises(ReviewError): + parse_csv(value) + self.assertEqual(len(parse_csv(b"sku,name,price_cents\n" + b"A,a,1\n" * 200)), 200) + + def test_invalid_values_survive_for_correction(self): + rows = parse_csv(b"sku,name,price_cents\nA,,3.50\nA,Second,01\n") + issues = validate_records(rows) + self.assertEqual(rows[0]["price_cents"], "3.50") + self.assertEqual(len(issues[0]), 3) + self.assertEqual(len(issues[1]), 2) + + def test_price_limits_and_unicode_controls(self): + for value in ["-1", "+1", "1.0", "1e3", "01", "1000000001", "9", ""]: + self.assertTrue( + validate_records([{"sku": "A", "name": "Item", "price_cents": value}])[0] + ) + self.assertEqual( + validate_records([{"sku": "A", "name": "Item", "price_cents": "1000000000"}]), [[]] + ) + for control in ["\x7f", "\x85", "\t"]: + self.assertTrue( + validate_records([{"sku": "A", "name": "Item" + control, "price_cents": "1"}])[0] + ) + + def test_formula_text_is_literal_and_catalogue_conflicts(self): + rows = parse_csv(b"sku,name,price_cents\nA,=1+1,1\n") + self.assertEqual(rows[0]["name"], "=1+1") + self.assertEqual(validate_records(rows), [[]]) + self.assertIn("already exists", validate_records(rows, ["A"])[0][0]) + + +if __name__ == "__main__": + unittest.main() From bfeeafbdeafa99b6669f38673754fcf9dda2529b Mon Sep 17 00:00:00 2001 From: sdairs Date: Fri, 2 Oct 2026 15:26:44 +0100 Subject: [PATCH 2/3] Update browser test helper to Playwright 1.56.0 --- applications/csv-review-desk/requirements-dev.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/applications/csv-review-desk/requirements-dev.txt b/applications/csv-review-desk/requirements-dev.txt index 30fb970d..cbff9cd7 100644 --- a/applications/csv-review-desk/requirements-dev.txt +++ b/applications/csv-review-desk/requirements-dev.txt @@ -20,7 +20,7 @@ numpy==2.5.3 packaging==26.3 pandas==3.0.6 pillow==12.3.0 -playwright==1.55.0 +playwright==1.56.0 protobuf==7.36.2 psycopg==3.3.6 psycopg-binary==3.3.6 From d312cd1a0e1ab8d09ed692ce5014bf5b6604b169 Mon Sep 17 00:00:00 2001 From: sdairs Date: Fri, 2 Oct 2026 21:31:27 +0100 Subject: [PATCH 3/3] docs: remove public beta qualifier --- applications/csv-review-desk/README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/applications/csv-review-desk/README.md b/applications/csv-review-desk/README.md index 0cfcc76e..ac60a7ad 100644 --- a/applications/csv-review-desk/README.md +++ b/applications/csv-review-desk/README.md @@ -1,6 +1,6 @@ # CSV review desk -Stage a CSV, correct its invalid rows, and approve the complete batch into a catalogue backed by [ClickHouse Managed Postgres (public beta)](https://clickhouse.com/docs/products/managed-postgres/overview). This small Streamlit workbench uses SQLAlchemy, psycopg, Alembic and pandas. Invalid values and their feedback remain durable, so closing the browser does not lose a review. +Stage a CSV, correct its invalid rows, and approve the complete batch into a catalogue backed by [ClickHouse Managed Postgres](https://clickhouse.com/docs/products/managed-postgres/overview). This small Streamlit workbench uses SQLAlchemy, psycopg, Alembic and pandas. Invalid values and their feedback remain durable, so closing the browser does not lose a review. Two views cover the workflow: upload/list saved batches, then review a batch in a fixed-row data editor. Approval inserts new catalogue items; it never updates an existing SKU. This is a trusted local workbench bound to loopback, without application authentication. Anyone who can reach the process has the same operator authority. Public deployment requires an authentication boundary and appropriate HTTPS/proxy configuration; it is outside this example's tested scope.