From 86c6d0ada5615e1cae1ab20f664ebcdf7a59264a Mon Sep 17 00:00:00 2001 From: sdairs Date: Fri, 2 Oct 2026 18:17:07 +0100 Subject: [PATCH 1/2] Add Scala experiment log with immutable retained requests --- .github/workflows/experiment-log.yml | 21 + applications/experiment-log/.env.example | 10 + applications/experiment-log/.gitignore | 11 + applications/experiment-log/.scalafmt.conf | 3 + applications/experiment-log/README.md | 124 ++++++ applications/experiment-log/build.sbt | 41 ++ .../experiment-log/dependency-lock.txt | 109 +++++ .../experiment-log/migrations/001-down.sql | 11 + .../experiment-log/migrations/001-up.sql | 116 ++++++ .../experiment-log/project/build.properties | 1 + .../experiment-log/project/plugins.sbt | 1 + .../experiment-log/project/repositories | 3 + applications/experiment-log/sample-run.json | 9 + .../experiment-log/scripts/bootstrap.sql | 12 + .../experiment-log/scripts/grants.sql | 4 + .../experiment-log/scripts/migrate.sh | 16 + applications/experiment-log/scripts/sbtw | 15 + applications/experiment-log/scripts/seed.sql | 5 + .../main/scala/experimentlog/Database.scala | 45 +++ .../src/main/scala/experimentlog/Domain.scala | 214 ++++++++++ .../src/main/scala/experimentlog/Main.scala | 130 ++++++ .../src/main/scala/experimentlog/Store.scala | 126 ++++++ .../scala/experimentlog/DomainSuite.scala | 60 +++ .../test/scala/experimentlog/LiveProbe.scala | 188 +++++++++ .../experiment-log/tests/acceptance.py | 372 ++++++++++++++++++ .../experiment-log/tests/requirements.in | 3 + .../experiment-log/tests/requirements.txt | 10 + 27 files changed, 1660 insertions(+) create mode 100644 .github/workflows/experiment-log.yml create mode 100644 applications/experiment-log/.env.example create mode 100644 applications/experiment-log/.gitignore create mode 100644 applications/experiment-log/.scalafmt.conf create mode 100644 applications/experiment-log/README.md create mode 100644 applications/experiment-log/build.sbt create mode 100644 applications/experiment-log/dependency-lock.txt create mode 100644 applications/experiment-log/migrations/001-down.sql create mode 100644 applications/experiment-log/migrations/001-up.sql create mode 100644 applications/experiment-log/project/build.properties create mode 100644 applications/experiment-log/project/plugins.sbt create mode 100644 applications/experiment-log/project/repositories create mode 100644 applications/experiment-log/sample-run.json create mode 100644 applications/experiment-log/scripts/bootstrap.sql create mode 100644 applications/experiment-log/scripts/grants.sql create mode 100755 applications/experiment-log/scripts/migrate.sh create mode 100755 applications/experiment-log/scripts/sbtw create mode 100644 applications/experiment-log/scripts/seed.sql create mode 100644 applications/experiment-log/src/main/scala/experimentlog/Database.scala create mode 100644 applications/experiment-log/src/main/scala/experimentlog/Domain.scala create mode 100644 applications/experiment-log/src/main/scala/experimentlog/Main.scala create mode 100644 applications/experiment-log/src/main/scala/experimentlog/Store.scala create mode 100644 applications/experiment-log/src/test/scala/experimentlog/DomainSuite.scala create mode 100644 applications/experiment-log/src/test/scala/experimentlog/LiveProbe.scala create mode 100644 applications/experiment-log/tests/acceptance.py create mode 100644 applications/experiment-log/tests/requirements.in create mode 100644 applications/experiment-log/tests/requirements.txt diff --git a/.github/workflows/experiment-log.yml b/.github/workflows/experiment-log.yml new file mode 100644 index 00000000..032356f5 --- /dev/null +++ b/.github/workflows/experiment-log.yml @@ -0,0 +1,21 @@ +name: Experiment log +on: + push: + paths: ['applications/experiment-log/**', '.github/workflows/experiment-log.yml'] + pull_request: + paths: ['applications/experiment-log/**', '.github/workflows/experiment-log.yml'] +permissions: + contents: read +jobs: + build: + runs-on: ubuntu-latest + defaults: + run: + working-directory: applications/experiment-log + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-java@v4 + with: + distribution: temurin + java-version: '21' + - run: scripts/sbtw clean scalafmtCheckAll compile test verifyDependencyLock diff --git a/applications/experiment-log/.env.example b/applications/experiment-log/.env.example new file mode 100644 index 00000000..11d1f248 --- /dev/null +++ b/applications/experiment-log/.env.example @@ -0,0 +1,10 @@ +# Copy values from your own dedicated Cloud Postgres receipt. +PGHOST=your-service-hostname +PGPORT=5432 +PGDATABASE=postgres +PGUSER=experiment_app +PGPASSWORD=replace-with-runtime-role-password +PGSSLROOTCERT=/absolute/path/to/cloud-ca.pem +PGSSLMODE=verify-full +PROJECT_A_TOKEN=replace-with-a-distinct-random-URL-safe-token-at-least-32-characters +PROJECT_B_TOKEN=replace-with-another-distinct-random-URL-safe-token-at-least-32-characters diff --git a/applications/experiment-log/.gitignore b/applications/experiment-log/.gitignore new file mode 100644 index 00000000..c3756942 --- /dev/null +++ b/applications/experiment-log/.gitignore @@ -0,0 +1,11 @@ +target/ +project/target/ +project/project/ +.local/ +.bsp/ +.metals/ +.venv/ +__pycache__/ +.env +*.env +*.pem diff --git a/applications/experiment-log/.scalafmt.conf b/applications/experiment-log/.scalafmt.conf new file mode 100644 index 00000000..5f191d11 --- /dev/null +++ b/applications/experiment-log/.scalafmt.conf @@ -0,0 +1,3 @@ +version = "3.11.5" +runner.dialect = scala3 +maxColumn = 100 diff --git a/applications/experiment-log/README.md b/applications/experiment-log/README.md new file mode 100644 index 00000000..5f639ffb --- /dev/null +++ b/applications/experiment-log/README.md @@ -0,0 +1,124 @@ +# Experiment log + +A small Scala API records supplied experiment results in ClickHouse Managed Postgres (public beta). Two seeded projects have separate bearer tokens. Each immutable run contains a title, a bounded flat JSONB configuration and 1–10 exact decimal measurements. These are synthetic examples, not model training or scientific validation. + +## Requirements and verified versions + +Use a native Linux environment with OpenJDK 21, curl, PostgreSQL `psql`, and a dedicated Cloud Postgres service. The checked set is Scala 3.9.0, sbt 1.13.0, http4s/Ember 0.23.38, Cats Effect 3.7.1, doobie **1.0.0-RC13** (a maintained published release candidate, not stable 1.0), Circe 0.14.16, HikariCP 7.1.0 and pgJDBC 42.7.13. Exact direct versions and the selected transitive dependency record are committed. The wrapper verifies its pinned sbt launcher SHA-256; `verifyDependencyLock` rejects resolver drift. No snapshots are used. + +The native CPU build was checked before allocating the database. Testing used Ubuntu 24.04 ARM64, OpenJDK 21.0.12.1 and PostgreSQL 18.6 on 2 October 2026. See [Cloud Postgres documentation](https://clickhouse.com/docs/products/managed-postgres/) for provisioning and the service CA. + +## Dedicated database setup + +Run from this application directory. Keep receipt credentials and the downloaded Cloud CA outside the clone, restrict private files to mode 600, and do not commit them. Never use this bootstrap on a shared database: it changes PUBLIC database/schema creation privileges. + +Install/configure [clickhousectl](https://github.com/ClickHouse/clickhousectl) with your own Cloud API credentials. Verify current supported size/region and [pricing](https://clickhouse.com/pricing) before creation. This dedicated fixture used AWS us-east-1, `c6gd.large`, PostgreSQL 18 and no HA; it incurs compute/storage charges. Stopping the API does not remove the database. + +```bash +mkdir -p $HOME/experiment-private +chmod 700 $HOME/experiment-private +clickhousectl cloud postgres create --org-id YOUR_ORG_ID \ + --name experiment-log --region us-east-1 --provider aws \ + --size c6gd.large --pg-version 18 --ha-type none \ + --tag purpose=experiment-log --json > $HOME/experiment-private/create.json +chmod 600 $HOME/experiment-private/create.json +clickhousectl cloud postgres get YOUR_SERVICE_ID --org-id YOUR_ORG_ID --json +``` + +Take `YOUR_SERVICE_ID` from the private receipt. Wait until `get` reports running, with its hostname populated, before retrieving the CA: + +```bash +clickhousectl cloud postgres certs get YOUR_SERVICE_ID \ + --org-id YOUR_ORG_ID --output $HOME/experiment-private/ca.pem +chmod 600 $HOME/experiment-private/ca.pem +``` + +Create four private shell environment files. Each includes the receipt hostname, port, database, `PGSSLROOTCERT` (absolute CA path), and `PGSSLMODE=verify-full`: + +* `admin.env`: receipt `PGUSER`/`PGPASSWORD`, plus `PG_MIGRATION_PASSWORD` and `PG_APP_PASSWORD` for two newly generated role passwords. +* `migration.env`: `PGUSER=experiment_owner` and its password. +* `runtime.env`: `PGUSER=experiment_app` and its password, plus distinct `PROJECT_A_TOKEN` and `PROJECT_B_TOKEN`. +* `test.env`: `TEST_OWNER_USER=experiment_owner` and `TEST_OWNER_PASSWORD` for opt-in acceptance only. + +Generate each token with `python3 -c 'import secrets; print(secrets.token_urlsafe(32))'`. Do not put tokens in URLs or screenshots. The API uses explicit Authorization headers; there are no cookies, browser sign-in, CORS allowances or UI. Bind locally for this tutorial; use HTTPS and proper token distribution for a real deployment. Tokens identify fixed seeded project UUIDs and are loaded at process startup. They are long-lived bearer credentials; this example does not implement per-user identity, expiry, audit trails or token revocation storage. + +In separate subshells, use the correct role for each step: + +```bash +(set -a; source $HOME/experiment-private/admin.env; set +a + psql -X -v ON_ERROR_STOP=1 -f scripts/bootstrap.sql) +(set -a; source $HOME/experiment-private/migration.env; set +a + scripts/migrate.sh up + psql -X -v ON_ERROR_STOP=1 -f scripts/grants.sql + psql -X -v ON_ERROR_STOP=1 -f scripts/seed.sql) +``` + +The bootstrap creates `experiment_owner`, `experiment_app` and the schema. Runtime can SELECT and INSERT only the required columns; it cannot update/delete runs or children, set timestamps/transaction IDs, or create tables. The shared runtime database role is trusted across both projects; project isolation is enforced by the API's server-derived identity, not PostgreSQL row-level security. + +`migrate.sh up` takes an advisory transaction lock and verifies the recorded file checksum on repeat. Seed is repeat-safe. `migrate.sh down` removes app data and tables; use it only on this dedicated fixture. To check a fresh migration cycle before recording data, run down/up as owner, then reapply grants and seed. Keep administrator credentials out of the runtime environment. + +## Build and run + +```bash +scripts/sbtw clean scalafmtCheckAll compile test verifyDependencyLock writeRuntimeClasspath +``` + +Start from a fresh shell containing **only** `runtime.env` and normal process variables. Do not source administrator, migration or test files in that shell: + +```bash +set -a; source $HOME/experiment-private/runtime.env; set +a +java -Xmx768m -XX:ActiveProcessorCount=2 \ + -cp "$(cat .local/runtime-classpath.txt)" experimentlog.Main +``` + +The server binds `127.0.0.1:8080` and checks that the two seeded project IDs match its configured token identities. pgJDBC requires the actual Cloud CA with `sslmode=verify-full`, including hostname verification. Hikari owns four connections, with bounded connect/validation waits. SQL statements are limited to 10 seconds, locks to 8, socket reads to 15. Headers have a 5-second deadline, request JSON is at most 16 KiB with a 5-second body deadline, and the response middleware has a 30-second deadline. SQL parameter logging is disabled. + +## Record, replay and search + +```bash +curl --fail-with-body http://127.0.0.1:8080/runs \ + -H "Authorization: Bearer $PROJECT_A_TOKEN" \ + -H 'Content-Type: application/json' --data-binary @sample-run.json +curl --fail-with-body http://127.0.0.1:8080/runs/search \ + -H "Authorization: Bearer $PROJECT_A_TOKEN" \ + -H 'Content-Type: application/json' \ + --data '{"config":{"optimizer":"adam"},"title":"comparison","limit":5}' +``` + +Registration returns 201; an identical retained request returns 200 and its original complete response. Reuse the sample UUID to observe replay, or generate a new UUID for another run. Reusing a UUID with changed semantic content returns 409. The UUID is retained for the lifetime of its project data, without expiry or deletion in this example. Whitespace around titles is trimmed; configuration object order, measurement array order and decimal trailing zeros do not change the semantic payload. Changing a configuration string's whitespace does change it. + +Measurement names match `[a-z][a-z0-9_]{0,31}` and are unique. Values are **plain decimal strings**, absolute value at most 1,000,000 and at most six fractional digits. They pass through Scala `BigDecimal`, doobie and unconstrained PostgreSQL `numeric`; a CHECK rejects excess scale instead of silently rounding it. Responses return decimal strings. No `Double` participates in the storage path. + +Configuration is an object of at most ten ASCII keys (`[A-Za-z][A-Za-z0-9_]{0,23}`), with bounded strings, booleans or exact JSON numbers of magnitude at most 1,000,000 and normalized scale at most six. Null, arrays, nested objects, controls, unpaired Unicode surrogates, duplicate object keys and oversized numeric/exponent lexemes are rejected. GET `/runs/:id` retrieves a project-scoped run; GET `/project` returns the authenticated project ID. + +Search applies parameterized JSONB `@>` containment plus an optional literal title substring. `%`, `_` and backslash are escaped. Empty `{}` matches every configuration. Page size is 1–20 (default 10), ordered by `(created_at,id)` descending. Pass `nextCursor` as `cursor` to continue. This is a validated opaque keyset cursor over a **live** listing, not a frozen snapshot: later inserts can appear above an existing cursor. Populated pages use two prepared SELECT statements, loading all selected child rows together; immutable committed runs make the pair consistent. The GIN `jsonb_path_ops` index supports containment; a small fixture may correctly use a sequential scan. No speed claim is made. + +The parent and complete measurement set are one `ConnectionIO` with one `transact` at the HTTP boundary. The unique `(project_id,request_id)` key arbitrates races; the losing `ON CONFLICT DO NOTHING` statement is followed by a fresh READ COMMITTED lookup. A deferred constraint checks complete children at commit, and a transaction-ID trigger prevents attaching later measurements to an already committed run. This database role is trusted to submit its own consistent initial payload; it is not a sandbox for arbitrary hostile SQL. + +## Verification + +CI runs only formatting, compile, six boundary/semantic unit tests and dependency-record verification, without database credentials. Opt-in live checks need the running API, runtime env and test env: + +```bash +python3 -m venv .venv +.venv/bin/pip install -r tests/requirements.txt +.venv/bin/pip check +(set -a; source $HOME/experiment-private/runtime.env; source $HOME/experiment-private/test.env; set +a + .venv/bin/python tests/acceptance.py) +scripts/sbtw Test/compile writeTestClasspath +(set -a; source $HOME/experiment-private/runtime.env; set +a + java -Xmx512m -cp "$(cat .local/test-classpath.txt)" experimentlog.LiveProbe) +``` + +Acceptance uses independent HTTP clients with observed database lock contention, a controlled owner-installed after-child failure and runtime privilege checks. The test-only TLS hostname control passes the actual TLS session to pgJDBC's native verifier with a substituted incorrect name; the application retains its default verifier. LiveProbe also performs actual doobie SQL type analysis, observes two prepared statements for a populated search, and prints the actual GIN definition and unforced plan. Its test classes are absent from the production classpath. Live checks are never silently skipped by CI. + +## Stop and remove the dedicated fixture + +Stop the JVM first. Delete only the exact Cloud service ID you created, using the explicit organization: + +```bash +clickhousectl cloud postgres delete YOUR_SERVICE_ID --org-id YOUR_ORG_ID --json +clickhousectl cloud postgres list --org-id YOUR_ORG_ID --json +``` + +Deletion is asynchronous; verify that exact ID becomes absent. If retaining the service, use the administrator receipt (not runtime credentials) to drop the dedicated schema/roles after stopping all connections. Never run cleanup against another service. The VM can be stopped and preserved; source and private evidence should be saved separately before shutdown. diff --git a/applications/experiment-log/build.sbt b/applications/experiment-log/build.sbt new file mode 100644 index 00000000..e953254e --- /dev/null +++ b/applications/experiment-log/build.sbt @@ -0,0 +1,41 @@ +ThisBuild / scalaVersion := "3.9.0" +ThisBuild / organization := "examples.clickhouse" +ThisBuild / version := "0.1.0" + +val http4sVersion = "0.23.38" +val doobieVersion = "1.0.0-RC13" // Maintained published release candidate, not stable 1.0. + +libraryDependencies ++= Seq( + "org.http4s" %% "http4s-ember-server" % http4sVersion, + "org.http4s" %% "http4s-dsl" % http4sVersion, + "org.http4s" %% "http4s-circe" % http4sVersion, + "org.typelevel" %% "cats-effect" % "3.7.1", + "org.typelevel" %% "doobie-core" % doobieVersion, + "org.typelevel" %% "doobie-hikari" % doobieVersion, + "org.typelevel" %% "doobie-postgres" % doobieVersion, + "org.typelevel" %% "doobie-postgres-circe" % doobieVersion, + "io.circe" %% "circe-core" % "0.14.16", + "io.circe" %% "circe-parser" % "0.14.16", + "com.zaxxer" % "HikariCP" % "7.1.0", + "org.postgresql" % "postgresql" % "42.7.13", + "org.slf4j" % "slf4j-simple" % "2.0.17", + "org.scalameta" %% "munit" % "1.3.6" % Test +) + +// Record selected modules from the actual resolver, including test dependencies. +val resolvedModules = taskKey[String]("Canonical selected dependency record") +resolvedModules := (Compile / update).value.configurations.flatMap(_.modules) + .filterNot(_.evicted).map(m => s"${m.module.organization}:${m.module.name}:${m.module.revision}") + .distinct.sorted.mkString("", "\n", "\n") +val writeDependencyLock = taskKey[Unit]("Write resolved dependency record intentionally") +writeDependencyLock := IO.write(baseDirectory.value / "dependency-lock.txt", resolvedModules.value) +val verifyDependencyLock = taskKey[Unit]("Fail when resolved dependencies differ from committed record") +verifyDependencyLock := { + val path = baseDirectory.value / "dependency-lock.txt" + require(path.exists && IO.read(path) == resolvedModules.value, "Dependency record differs; review before writeDependencyLock") +} +val writeRuntimeClasspath = taskKey[Unit]("Write native production JVM classpath") +writeRuntimeClasspath := IO.write(baseDirectory.value / ".local/runtime-classpath.txt", (Runtime / fullClasspath).value.files.mkString(java.io.File.pathSeparator)) + +val writeTestClasspath = taskKey[Unit]("Write native live-check JVM classpath") +writeTestClasspath := IO.write(baseDirectory.value / ".local/test-classpath.txt", (Test / fullClasspath).value.files.mkString(java.io.File.pathSeparator)) diff --git a/applications/experiment-log/dependency-lock.txt b/applications/experiment-log/dependency-lock.txt new file mode 100644 index 00000000..fcba0d14 --- /dev/null +++ b/applications/experiment-log/dependency-lock.txt @@ -0,0 +1,109 @@ +co.fs2:fs2-core_3:3.14.0 +co.fs2:fs2-io_3:3.14.0 +com.comcast:ip4s-core_3:3.8.0 +com.fasterxml.jackson.core:jackson-annotations:2.21 +com.fasterxml.jackson.core:jackson-core:2.13.4 +com.fasterxml.jackson.core:jackson-databind:2.13.4.2 +com.fasterxml.jackson.datatype:jackson-datatype-jsr310:2.13.2 +com.twitter:hpack:1.0.2 +com.vladsch.flexmark:flexmark-ext-anchorlink:0.64.8 +com.vladsch.flexmark:flexmark-ext-autolink:0.64.8 +com.vladsch.flexmark:flexmark-ext-emoji:0.64.8 +com.vladsch.flexmark:flexmark-ext-gfm-strikethrough:0.64.8 +com.vladsch.flexmark:flexmark-ext-gfm-tasklist:0.64.8 +com.vladsch.flexmark:flexmark-ext-ins:0.64.8 +com.vladsch.flexmark:flexmark-ext-superscript:0.64.8 +com.vladsch.flexmark:flexmark-ext-tables:0.64.8 +com.vladsch.flexmark:flexmark-ext-wikilink:0.64.8 +com.vladsch.flexmark:flexmark-ext-yaml-front-matter:0.64.8 +com.vladsch.flexmark:flexmark-jira-converter:0.64.8 +com.vladsch.flexmark:flexmark-util-ast:0.64.8 +com.vladsch.flexmark:flexmark-util-builder:0.64.8 +com.vladsch.flexmark:flexmark-util-collection:0.64.8 +com.vladsch.flexmark:flexmark-util-data:0.64.8 +com.vladsch.flexmark:flexmark-util-dependency:0.64.8 +com.vladsch.flexmark:flexmark-util-format:0.64.8 +com.vladsch.flexmark:flexmark-util-html:0.64.8 +com.vladsch.flexmark:flexmark-util-misc:0.64.8 +com.vladsch.flexmark:flexmark-util-options:0.64.8 +com.vladsch.flexmark:flexmark-util-sequence:0.64.8 +com.vladsch.flexmark:flexmark-util-visitor:0.64.8 +com.vladsch.flexmark:flexmark-util:0.64.8 +com.vladsch.flexmark:flexmark:0.64.8 +com.zaxxer:HikariCP:7.1.0 +io.circe:circe-core_3:0.14.16 +io.circe:circe-jawn_3:0.14.16 +io.circe:circe-numbers_3:0.14.16 +io.circe:circe-parser_3:0.14.16 +io.get-coursier:interface:1.0.29-M4 +junit:junit:4.13.2 +nl.big-o:liqp:0.9.2.3 +org.checkerframework:checker-qual:3.55.1 +org.hamcrest:hamcrest-core:1.3 +org.http4s:http4s-circe_3:0.23.38 +org.http4s:http4s-core_3:0.23.38 +org.http4s:http4s-crypto_3:0.2.5 +org.http4s:http4s-dsl_3:0.23.38 +org.http4s:http4s-ember-core_3:0.23.38 +org.http4s:http4s-ember-server_3:0.23.38 +org.http4s:http4s-jawn_3:0.23.38 +org.http4s:http4s-server_3:0.23.38 +org.jetbrains:annotations:24.0.1 +org.jline:jline-native:4.0.14 +org.jline:jline-reader:4.0.14 +org.jline:jline-terminal-jni:4.0.14 +org.jline:jline-terminal:4.0.14 +org.jsoup:jsoup:1.22.2 +org.log4s:log4s_3:1.10.0 +org.nibor.autolink:autolink:0.6.0 +org.portable-scala:portable-scala-reflect_2.13:1.1.3 +org.postgresql:postgresql:42.7.13 +org.scala-lang.modules:scala-asm:9.9.0-scala-1 +org.scala-lang:scala-library:3.9.0 +org.scala-lang:scala3-compiler_3:3.9.0 +org.scala-lang:scala3-directives-parser_3:3.9.0 +org.scala-lang:scala3-interfaces:3.9.0 +org.scala-lang:scala3-library_3:3.9.0 +org.scala-lang:scala3-repl_3:3.9.0 +org.scala-lang:scala3-tasty-inspector_3:3.9.0 +org.scala-lang:scaladoc_3:3.9.0 +org.scala-lang:tasty-core_3:3.9.0 +org.scala-sbt:compiler-interface:1.12.0 +org.scala-sbt:test-interface:1.0 +org.scala-sbt:util-interface:1.11.5 +org.scalameta:junit-interface:1.3.6 +org.scalameta:munit-diff_3:1.3.6 +org.scalameta:munit_3:1.3.6 +org.scodec:scodec-bits_3:1.2.5 +org.slf4j:slf4j-api:1.7.36 +org.slf4j:slf4j-api:2.0.17 +org.slf4j:slf4j-simple:2.0.17 +org.snakeyaml:snakeyaml-engine:3.0.1 +org.tpolecat:typename_3:1.1.2 +org.typelevel:algebra_3:2.13.0 +org.typelevel:case-insensitive_3:1.5.0 +org.typelevel:cats-collections-core_3:0.9.10 +org.typelevel:cats-core_3:2.13.0 +org.typelevel:cats-effect-kernel_3:3.7.1 +org.typelevel:cats-effect-std_3:3.7.1 +org.typelevel:cats-effect_3:3.7.1 +org.typelevel:cats-free_3:2.13.0 +org.typelevel:cats-kernel_3:2.13.0 +org.typelevel:cats-mtl_3:1.7.0 +org.typelevel:cats-parse_3:1.1.0 +org.typelevel:doobie-core_3:1.0.0-RC13 +org.typelevel:doobie-free_3:1.0.0-RC13 +org.typelevel:doobie-hikari_3:1.0.0-RC13 +org.typelevel:doobie-postgres-circe_3:1.0.0-RC13 +org.typelevel:doobie-postgres_3:1.0.0-RC13 +org.typelevel:idna4s-core_3:0.1.0 +org.typelevel:jawn-fs2_3:2.6.0 +org.typelevel:jawn-parser_3:1.7.0 +org.typelevel:literally_3:1.2.0 +org.typelevel:log4cats-core_3:2.8.0 +org.typelevel:log4cats-slf4j_3:2.8.0 +org.typelevel:vault_3:3.7.0 +tools.jackson.core:jackson-core:3.1.2 +tools.jackson.core:jackson-databind:3.1.2 +tools.jackson.dataformat:jackson-dataformat-yaml:3.1.2 +ua.co.k:strftime4j:1.0.6 diff --git a/applications/experiment-log/migrations/001-down.sql b/applications/experiment-log/migrations/001-down.sql new file mode 100644 index 00000000..93dea2a4 --- /dev/null +++ b/applications/experiment-log/migrations/001-down.sql @@ -0,0 +1,11 @@ +BEGIN; +SELECT pg_advisory_xact_lock(hashtextextended('experiment_log:001', 0)); +DROP TABLE experiment_log.measurements; +DROP TABLE experiment_log.runs; +DROP TABLE experiment_log.projects; +DROP FUNCTION experiment_log.complete_run(); +DROP FUNCTION experiment_log.initial_measurement(); +DROP FUNCTION experiment_log.valid_measurement_payload(jsonb); +DROP FUNCTION experiment_log.valid_config(jsonb); +DELETE FROM experiment_log.schema_migrations WHERE version=1; +COMMIT; diff --git a/applications/experiment-log/migrations/001-up.sql b/applications/experiment-log/migrations/001-up.sql new file mode 100644 index 00000000..adaae369 --- /dev/null +++ b/applications/experiment-log/migrations/001-up.sql @@ -0,0 +1,116 @@ +\getenv migration_sha MIGRATION_SHA +BEGIN; +SELECT pg_advisory_xact_lock(hashtextextended('experiment_log:001', 0)); +CREATE TABLE IF NOT EXISTS experiment_log.schema_migrations ( + version integer PRIMARY KEY, + checksum text NOT NULL, + applied_at timestamptz NOT NULL DEFAULT now() +); +SELECT EXISTS (SELECT 1 FROM experiment_log.schema_migrations WHERE version=1) AS applied \gset +\if :applied + SELECT checksum = :'migration_sha' AS same FROM experiment_log.schema_migrations WHERE version=1 \gset + \if :same + \echo Migration 001 already applied with matching checksum. + \else + \echo Migration checksum differs; refusing to continue. + \quit 3 + \endif +\else +CREATE FUNCTION experiment_log.valid_config(settings jsonb) RETURNS boolean +LANGUAGE plpgsql IMMUTABLE AS $$ +DECLARE entry record; number numeric; +BEGIN + IF jsonb_typeof(settings) IS DISTINCT FROM 'object' THEN RETURN false; END IF; + IF (SELECT count(*) FROM jsonb_object_keys(settings)) > 10 THEN RETURN false; END IF; + FOR entry IN SELECT key,value FROM jsonb_each(settings) LOOP + IF entry.key !~ '^[A-Za-z][A-Za-z0-9_]{0,23}$' THEN RETURN false; END IF; + CASE jsonb_typeof(entry.value) + WHEN 'string' THEN IF length(entry.value #>> '{}') > 120 THEN RETURN false; END IF; + WHEN 'boolean' THEN NULL; + WHEN 'number' THEN + number := (entry.value #>> '{}')::numeric; + IF abs(number) > 1000000 OR round(number,6) <> number THEN RETURN false; END IF; + ELSE RETURN false; + END CASE; + END LOOP; + RETURN true; +END $$; +CREATE FUNCTION experiment_log.valid_measurement_payload(values_json jsonb) RETURNS boolean +LANGUAGE plpgsql IMMUTABLE AS $$ +DECLARE entry record; value_text text; total integer; +BEGIN + IF jsonb_typeof(values_json) IS DISTINCT FROM 'object' THEN RETURN false; END IF; + SELECT count(*) INTO total FROM jsonb_object_keys(values_json); + IF total NOT BETWEEN 1 AND 10 THEN RETURN false; END IF; + FOR entry IN SELECT key,value FROM jsonb_each(values_json) LOOP + IF entry.key !~ '^[a-z][a-z0-9_]{0,31}$' OR jsonb_typeof(entry.value) <> 'string' THEN RETURN false; END IF; + value_text := entry.value #>> '{}'; + IF length(value_text) > 16 OR value_text !~ '^-?(0|[1-9][0-9]{0,6})(\.[0-9]{1,6})?$' THEN RETURN false; END IF; + IF abs(value_text::numeric) > 1000000 THEN RETURN false; END IF; + END LOOP; + RETURN true; +END $$; +CREATE TABLE experiment_log.projects ( + id uuid PRIMARY KEY, + name text NOT NULL CHECK (length(name) BETWEEN 1 AND 80) +); +CREATE TABLE experiment_log.runs ( + id uuid PRIMARY KEY, + project_id uuid NOT NULL REFERENCES experiment_log.projects(id), + request_id uuid NOT NULL, + title text NOT NULL CHECK (length(btrim(title)) BETWEEN 1 AND 120), + config jsonb NOT NULL CHECK (experiment_log.valid_config(config)), + request_payload jsonb NOT NULL CHECK ( + jsonb_typeof(request_payload) = 'object' + AND request_payload ?& ARRAY['title','config','measurements'] + AND jsonb_typeof(request_payload->'title') = 'string' + AND request_payload->>'title' = title AND request_payload->'config' = config + AND experiment_log.valid_measurement_payload(request_payload->'measurements') + AND pg_column_size(request_payload) <= 32768 + ), + created_at timestamptz NOT NULL DEFAULT now(), + created_xid xid8 NOT NULL DEFAULT pg_current_xact_id(), + UNIQUE(project_id,request_id), + UNIQUE(project_id,id) +); +CREATE TABLE experiment_log.measurements ( + project_id uuid NOT NULL, + run_id uuid NOT NULL, + name text NOT NULL CHECK (name ~ '^[a-z][a-z0-9_]{0,31}$'), + -- No typmod rounding: reject excess scale before any precision can be lost. + value numeric NOT NULL CHECK (abs(value) <= 1000000 AND scale(value) <= 6), + PRIMARY KEY(project_id,run_id,name), + FOREIGN KEY(project_id,run_id) REFERENCES experiment_log.runs(project_id,id) +); +CREATE INDEX runs_config_gin ON experiment_log.runs USING gin(config jsonb_path_ops); +CREATE INDEX runs_project_page ON experiment_log.runs(project_id,created_at DESC,id DESC); +CREATE FUNCTION experiment_log.initial_measurement() RETURNS trigger LANGUAGE plpgsql AS $$ +DECLARE parent_xid xid8; expected text; +BEGIN + SELECT created_xid,request_payload->'measurements'->>NEW.name INTO parent_xid,expected + FROM experiment_log.runs WHERE project_id=NEW.project_id AND id=NEW.run_id; + IF parent_xid IS DISTINCT FROM pg_current_xact_id() THEN + RAISE EXCEPTION 'Measurements belong only to the initial run transaction' USING ERRCODE='23514'; + END IF; + IF expected IS NULL OR NEW.value <> expected::numeric THEN + RAISE EXCEPTION 'Measurement does not match retained payload' USING ERRCODE='23514'; + END IF; + RETURN NEW; +END $$; +CREATE TRIGGER initial_measurement BEFORE INSERT ON experiment_log.measurements +FOR EACH ROW EXECUTE FUNCTION experiment_log.initial_measurement(); +CREATE FUNCTION experiment_log.complete_run() RETURNS trigger LANGUAGE plpgsql AS $$ +DECLARE actual integer; expected integer; +BEGIN + SELECT count(*) INTO expected FROM jsonb_object_keys(NEW.request_payload->'measurements'); + SELECT count(*) INTO actual FROM experiment_log.measurements WHERE project_id=NEW.project_id AND run_id=NEW.id; + IF actual <> expected OR actual NOT BETWEEN 1 AND 10 THEN + RAISE EXCEPTION 'Run requires its complete measurement set' USING ERRCODE='23514'; + END IF; + RETURN NEW; +END $$; +CREATE CONSTRAINT TRIGGER complete_run AFTER INSERT ON experiment_log.runs +DEFERRABLE INITIALLY DEFERRED FOR EACH ROW EXECUTE FUNCTION experiment_log.complete_run(); +INSERT INTO experiment_log.schema_migrations(version,checksum) VALUES(1,:'migration_sha'); +\endif +COMMIT; diff --git a/applications/experiment-log/project/build.properties b/applications/experiment-log/project/build.properties new file mode 100644 index 00000000..e0a1aa02 --- /dev/null +++ b/applications/experiment-log/project/build.properties @@ -0,0 +1 @@ +sbt.version=1.13.0 diff --git a/applications/experiment-log/project/plugins.sbt b/applications/experiment-log/project/plugins.sbt new file mode 100644 index 00000000..c0372223 --- /dev/null +++ b/applications/experiment-log/project/plugins.sbt @@ -0,0 +1 @@ +addSbtPlugin("org.scalameta" % "sbt-scalafmt" % "2.6.2") diff --git a/applications/experiment-log/project/repositories b/applications/experiment-log/project/repositories new file mode 100644 index 00000000..60add579 --- /dev/null +++ b/applications/experiment-log/project/repositories @@ -0,0 +1,3 @@ +[repositories] +local +maven-central: https://repo.maven.apache.org/maven2 diff --git a/applications/experiment-log/sample-run.json b/applications/experiment-log/sample-run.json new file mode 100644 index 00000000..d9fbec54 --- /dev/null +++ b/applications/experiment-log/sample-run.json @@ -0,0 +1,9 @@ +{ + "requestId": "ae9387cb-13e5-4e92-a1ef-ea61df9c8668", + "title": "Synthetic optimizer comparison", + "config": {"optimizer": "adam", "enabled": true, "rate": 0.000001}, + "measurements": [ + {"name": "accuracy", "value": "0.123456"}, + {"name": "loss", "value": "-0.000001"} + ] +} diff --git a/applications/experiment-log/scripts/bootstrap.sql b/applications/experiment-log/scripts/bootstrap.sql new file mode 100644 index 00000000..fc734b56 --- /dev/null +++ b/applications/experiment-log/scripts/bootstrap.sql @@ -0,0 +1,12 @@ +-- Administrator only; dedicated service. +\getenv owner_password PG_MIGRATION_PASSWORD +\getenv app_password PG_APP_PASSWORD +CREATE ROLE experiment_owner LOGIN PASSWORD :'owner_password'; +CREATE ROLE experiment_app LOGIN PASSWORD :'app_password'; +CREATE SCHEMA experiment_log AUTHORIZATION experiment_owner; +REVOKE ALL ON SCHEMA experiment_log FROM PUBLIC; +REVOKE CREATE ON SCHEMA public FROM PUBLIC; +REVOKE CREATE, TEMP ON DATABASE postgres FROM PUBLIC; +GRANT CONNECT ON DATABASE postgres TO experiment_owner, experiment_app; +ALTER ROLE experiment_owner SET search_path TO experiment_log, public; +ALTER ROLE experiment_app SET search_path TO experiment_log, public; diff --git a/applications/experiment-log/scripts/grants.sql b/applications/experiment-log/scripts/grants.sql new file mode 100644 index 00000000..c32e7ae3 --- /dev/null +++ b/applications/experiment-log/scripts/grants.sql @@ -0,0 +1,4 @@ +GRANT USAGE ON SCHEMA experiment_log TO experiment_app; +GRANT SELECT ON experiment_log.projects,experiment_log.runs,experiment_log.measurements TO experiment_app; +GRANT INSERT(id,project_id,request_id,title,config,request_payload) ON experiment_log.runs TO experiment_app; +GRANT INSERT(project_id,run_id,name,value) ON experiment_log.measurements TO experiment_app; diff --git a/applications/experiment-log/scripts/migrate.sh b/applications/experiment-log/scripts/migrate.sh new file mode 100755 index 00000000..4ef7c798 --- /dev/null +++ b/applications/experiment-log/scripts/migrate.sh @@ -0,0 +1,16 @@ +#!/usr/bin/env bash +set -euo pipefail +cd "$(dirname "$0")/.." +: "${PGSSLROOTCERT:?Cloud CA is required}" +: "${PGUSER:?Schema-owner role is required}" +test "$PGUSER" = experiment_owner +export PGSSLMODE=verify-full +case "${1:-up}" in + up) + export MIGRATION_SHA + MIGRATION_SHA=$(sha256sum migrations/001-up.sql | cut -d ' ' -f1) + psql -X -v ON_ERROR_STOP=1 -f migrations/001-up.sql + ;; + down) psql -X -v ON_ERROR_STOP=1 -f migrations/001-down.sql ;; + *) echo 'Use migrate.sh up or down' >&2; exit 2 ;; +esac diff --git a/applications/experiment-log/scripts/sbtw b/applications/experiment-log/scripts/sbtw new file mode 100755 index 00000000..43787e59 --- /dev/null +++ b/applications/experiment-log/scripts/sbtw @@ -0,0 +1,15 @@ +#!/usr/bin/env bash +set -euo pipefail +cd "$(dirname "$0")/.." +version=1.13.0 +expected=18d4a733eae2f70fd4d71a46e9da21d2eb7c9bf6ee85a67be9a9800e7c01403e +launcher="$PWD/.local/sbt/sbt-launch.jar" +mkdir -p "$(dirname "$launcher")" +if [ ! -f "$launcher" ]; then + curl --fail --location --retry 2 --connect-timeout 10 --max-time 120 \ + "https://repo.maven.apache.org/maven2/org/scala-sbt/sbt-launch/$version/sbt-launch-$version.jar" -o "$launcher" +fi +printf '%s %s\n' "$expected" "$launcher" | sha256sum --check --status +exec java -Xmx1024m -Xss2m -XX:ActiveProcessorCount=2 \ + -Dsbt.override.build.repos=true -Dsbt.repository.config="$PWD/project/repositories" \ + -Dsbt.supershell=false -jar "$launcher" "$@" diff --git a/applications/experiment-log/scripts/seed.sql b/applications/experiment-log/scripts/seed.sql new file mode 100644 index 00000000..59da0235 --- /dev/null +++ b/applications/experiment-log/scripts/seed.sql @@ -0,0 +1,5 @@ +-- Owner only; no tokens or scientific results stored in this seed. +INSERT INTO experiment_log.projects(id,name) VALUES +('00000000-0000-4000-8000-000000000001','Synthetic Alpha'), +('00000000-0000-4000-8000-000000000002','Synthetic Beta') +ON CONFLICT(id) DO NOTHING; diff --git a/applications/experiment-log/src/main/scala/experimentlog/Database.scala b/applications/experiment-log/src/main/scala/experimentlog/Database.scala new file mode 100644 index 00000000..42161fee --- /dev/null +++ b/applications/experiment-log/src/main/scala/experimentlog/Database.scala @@ -0,0 +1,45 @@ +package experimentlog + +import cats.effect.{IO, Resource} +import com.zaxxer.hikari.HikariConfig +import org.typelevel.doobie.hikari.HikariTransactor +import java.nio.file.{Files, Path} + +object Database { + def resource: Resource[IO, HikariTransactor[IO]] = + Resource + .eval(IO { + def required(key: String) = + sys.env.getOrElse(key, throw new IllegalArgumentException(s"$key is required")) + val ca = required("PGSSLROOTCERT") + require(Files.isRegularFile(Path.of(ca)), "Cloud CA file is required") + val host = required("PGHOST") + require(host.matches("[A-Za-z0-9.-]+"), "Invalid database hostname") + val database = sys.env.getOrElse("PGDATABASE", "postgres") + require(database.matches("[A-Za-z0-9_]+"), "Invalid database name") + val port = sys.env.getOrElse("PGPORT", "5432").toInt + require(port >= 1 && port <= 65535, "Invalid database port") + val config = new HikariConfig() + config.setDriverClassName("org.postgresql.Driver") + config.setJdbcUrl(s"jdbc:postgresql://$host:$port/$database") + config.setUsername(required("PGUSER")) + config.setPassword(required("PGPASSWORD")) + config.setMaximumPoolSize(4) + config.setMinimumIdle(0) + config.setConnectionTimeout(5000) + config.setValidationTimeout(2000) + config.setInitializationFailTimeout(10000) + config.setTransactionIsolation("TRANSACTION_READ_COMMITTED") + config.addDataSourceProperty("sslmode", "verify-full") + config.addDataSourceProperty("sslrootcert", ca) + config.addDataSourceProperty("connectTimeout", "10") + config.addDataSourceProperty("socketTimeout", "15") + config.addDataSourceProperty("ApplicationName", "experiment-log") + config.addDataSourceProperty( + "options", + "-csearch_path=experiment_log,public -cstatement_timeout=10000 -clock_timeout=8000" + ) + config + }) + .flatMap(config => HikariTransactor.fromHikariConfig[IO](config, logHandler = None)) +} diff --git a/applications/experiment-log/src/main/scala/experimentlog/Domain.scala b/applications/experiment-log/src/main/scala/experimentlog/Domain.scala new file mode 100644 index 00000000..492e331a --- /dev/null +++ b/applications/experiment-log/src/main/scala/experimentlog/Domain.scala @@ -0,0 +1,214 @@ +package experimentlog + +import java.nio.charset.{CodingErrorAction, StandardCharsets} +import java.time.Instant +import java.util.{Base64, UUID} +import io.circe.{Json, JsonObject} +import io.circe.jawn.JawnParser +import scala.util.Try + +final case class Rejection(message: String, status: Int = 422) extends RuntimeException(message) +final case class Measurement(name: String, value: BigDecimal) +final case class RunInput( + requestId: UUID, + title: String, + config: Json, + measurements: Vector[Measurement] +) { + val semantic: Json = Json.obj( + "title" -> Json.fromString(title), + "config" -> config, + "measurements" -> Json.fromFields( + measurements.sortBy(_.name).map(m => m.name -> Json.fromString(Domain.decimalText(m.value))) + ) + ) +} +final case class Cursor(createdAt: Instant, id: UUID) { + def encoded: String = Base64.getUrlEncoder.withoutPadding.encodeToString( + (createdAt.toString + "|" + id.toString).getBytes(StandardCharsets.US_ASCII) + ) +} +final case class Search(config: Json, title: String, limit: Int, cursor: Option[Cursor]) + +object Domain { + val MaxBody = 16384 + private val parser = new JawnParser(maxValueSize = Some(MaxBody), allowDuplicateKeys = false) + private val keyPattern = "[A-Za-z][A-Za-z0-9_]{0,23}".r + private val namePattern = "[a-z][a-z0-9_]{0,31}".r + private val decimalPattern = "-?(0|[1-9][0-9]{0,6})(\\.[0-9]{1,6})?".r + private val numberPattern = "-?(0|[1-9][0-9]*)(\\.[0-9]+)?([eE][+-]?[0-9]+)?".r + + def ensure(condition: Boolean, message: String): Unit = if (!condition) throw Rejection(message) + def decimalText(value: BigDecimal): String = value.bigDecimal.stripTrailingZeros.toPlainString + + def safeText(value: String, max: Int): String = { + ensure(value.codePointCount(0, value.length) <= max, "Text exceeds its character bound.") + var i = 0 + while (i < value.length) { + val c = value.charAt(i) + if (Character.isHighSurrogate(c)) { + ensure( + i + 1 < value.length && Character.isLowSurrogate(value.charAt(i + 1)), + "Invalid Unicode." + ) + i += 2 + } else { + ensure( + !Character.isLowSurrogate(c) && !Character.isISOControl(c), + "Controls or invalid Unicode are not allowed." + ) + i += 1 + } + } + value + } + + // Bound nesting and numeric lexemes before Circe/BigDecimal can expand an exponent. + private def lexicalBounds(raw: String): Unit = { + var i = 0 + var depth = 0 + var quoted = false + var escaped = false + while (i < raw.length) { + val c = raw.charAt(i) + if (quoted) { + if (escaped) escaped = false + else if (c == '\\') escaped = true + else if (c == '"') quoted = false + } else if (c == '"') quoted = true + else if (c == '{' || c == '[') { + depth += 1 + ensure(depth <= 4, "JSON nesting exceeds the envelope bound.") + } else if (c == '}' || c == ']') depth -= 1 + else if (c == '-' || c.isDigit) { + val start = i + while (i < raw.length && "0123456789.eE+-".contains(raw.charAt(i))) i += 1 + val token = raw.substring(start, i) + ensure( + token.length <= 40 && numberPattern.matches(token), + "Invalid or oversized numeric literal." + ) + val exponent = token.indexWhere(ch => ch == 'e' || ch == 'E') + if (exponent >= 0) { + val suffix = token.substring(exponent + 1) + ensure( + suffix.length <= 3 && Try(suffix.toInt).toOption.exists(n => math.abs(n) <= 12), + "Numeric exponent is out of bounds." + ) + } + i -= 1 + } + i += 1 + } + } + + def parse(bytes: Array[Byte]): Json = { + ensure(bytes.length <= MaxBody, "Request exceeds 16 KiB.") + val decoder = StandardCharsets.UTF_8.newDecoder + .onMalformedInput(CodingErrorAction.REPORT) + .onUnmappableCharacter(CodingErrorAction.REPORT) + val raw = Try(decoder.decode(java.nio.ByteBuffer.wrap(bytes)).toString) + .getOrElse(throw Rejection("Request must be valid UTF-8.", 400)) + lexicalBounds(raw) + parser.parse(raw).getOrElse(throw Rejection("Malformed JSON or duplicate object keys.", 400)) + } + + private def obj(json: Json, allowed: Set[String]): JsonObject = { + val fields = json.asObject.getOrElse(throw Rejection("An object is required.")) + ensure(fields.keys.forall(allowed), "Unknown field.") + fields + } + private def string(fields: JsonObject, key: String): String = + fields(key).flatMap(_.asString).getOrElse(throw Rejection(s"$key must be a string.")) + def uuid(value: String): UUID = { + ensure( + value.matches("[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}"), + "A canonical UUID shape is required." + ) + UUID.fromString(value) + } + def config(json: Json): Json = { + val fields = json.asObject.getOrElse(throw Rejection("Configuration must be an object.")) + ensure(fields.size <= 10, "Configuration permits at most 10 keys.") + Json.fromFields(fields.toVector.sortBy(_._1).map { case (key, value) => + ensure(keyPattern.matches(key), "Configuration key must match [A-Za-z][A-Za-z0-9_]{0,23}.") + val normalized = value.asString match { + case Some(text) => Json.fromString(safeText(text, 120)) + case None if value.isBoolean => value + case None if value.isNumber => + val decimal = value.asNumber + .flatMap(_.toBigDecimal) + .getOrElse(throw Rejection("Invalid decimal configuration value.")) + val stripped = decimal.bigDecimal.stripTrailingZeros + ensure( + decimal.abs <= BigDecimal(1000000) && math + .max(0, stripped.scale) <= 6 && stripped.precision <= 13, + "Configuration decimal exceeds exact range/scale." + ) + Json.fromBigDecimal(BigDecimal(stripped)) + case _ => + throw Rejection( + "Configuration values must be scalar string, boolean or number; null is not allowed." + ) + } + key -> normalized + }) + } + def run(json: Json): RunInput = { + val fields = obj(json, Set("requestId", "title", "config", "measurements")) + val request = uuid(string(fields, "requestId")) + val title = safeText(string(fields, "title"), 120).trim + ensure(title.nonEmpty, "Title is required.") + val settings = config(fields("config").getOrElse(throw Rejection("config is required."))) + val rows = fields("measurements") + .flatMap(_.asArray) + .getOrElse(throw Rejection("measurements must be an array.")) + ensure(rows.nonEmpty && rows.size <= 10, "Supply 1–10 measurements.") + val measurements = rows.map { json => + val values = obj(json, Set("name", "value")) + val name = string(values, "name") + ensure(namePattern.matches(name), "Measurement name must match [a-z][a-z0-9_]{0,31}.") + val raw = string(values, "value") + ensure( + raw.length <= 16 && decimalPattern.matches(raw), + "Measurement value must be a plain decimal string with at most 6 fractional digits." + ) + val decimal = BigDecimal(raw) + ensure(decimal.abs <= BigDecimal(1000000), "Measurement magnitude exceeds 1,000,000.") + Measurement(name, BigDecimal(decimal.bigDecimal.stripTrailingZeros)) + } + ensure( + measurements.map(_.name).distinct.size == measurements.size, + "Measurement names must be unique." + ) + RunInput(request, title, settings, measurements.sortBy(_.name)) + } + def cursor(value: String): Cursor = { + ensure(value.length <= 120 && value.matches("[A-Za-z0-9_-]+"), "Invalid cursor.") + val raw = Try(new String(Base64.getUrlDecoder.decode(value), StandardCharsets.US_ASCII)) + .getOrElse(throw Rejection("Invalid cursor.")) + val parts = raw.split("\\|", -1) + ensure(parts.length == 2, "Invalid cursor.") + val at = Try(Instant.parse(parts(0))).getOrElse(throw Rejection("Invalid cursor timestamp.")) + ensure( + at.isAfter(Instant.parse("1970-01-01T00:00:00Z")) && at.isBefore( + Instant.parse("2100-01-01T00:00:00Z") + ), + "Cursor timestamp out of bounds." + ) + val result = Cursor(at, uuid(parts(1))) + ensure(result.encoded == value, "Noncanonical cursor.") + result + } + def search(json: Json): Search = { + val fields = obj(json, Set("config", "title", "limit", "cursor")) + val settings = config(fields("config").getOrElse(Json.obj())) + val title = fields("title").map(_ => safeText(string(fields, "title"), 80)).getOrElse("") + val limit = fields("limit") + .map(_.asNumber.flatMap(_.toInt).getOrElse(throw Rejection("limit must be an integer."))) + .getOrElse(10) + ensure(limit >= 1 && limit <= 20, "limit must be 1–20.") + val after = fields("cursor").map(_ => cursor(string(fields, "cursor"))) + Search(settings, title, limit, after) + } +} diff --git a/applications/experiment-log/src/main/scala/experimentlog/Main.scala b/applications/experiment-log/src/main/scala/experimentlog/Main.scala new file mode 100644 index 00000000..533ac78e --- /dev/null +++ b/applications/experiment-log/src/main/scala/experimentlog/Main.scala @@ -0,0 +1,130 @@ +package experimentlog + +import cats.effect.{IO, IOApp} +import cats.syntax.all.* +import com.comcast.ip4s.* +import org.http4s.* +import org.http4s.circe.* +import org.http4s.dsl.io.* +import org.http4s.ember.server.EmberServerBuilder +import org.http4s.server.middleware.Timeout +import org.typelevel.doobie.implicits.* +import org.typelevel.doobie.* +import org.typelevel.doobie.postgres.implicits.* +import org.typelevel.doobie.Transactor +import io.circe.Json +import java.security.MessageDigest +import java.nio.charset.StandardCharsets +import java.util.UUID +import scala.concurrent.duration.* + +object Auth { + val projectA = UUID.fromString("00000000-0000-4000-8000-000000000001") + val projectB = UUID.fromString("00000000-0000-4000-8000-000000000002") + def configured: Vector[(Array[Byte], UUID)] = { + val values = Vector("PROJECT_A_TOKEN" -> projectA, "PROJECT_B_TOKEN" -> projectB).map { + case (key, id) => + val token = sys.env.getOrElse(key, throw new IllegalArgumentException(s"$key is required")) + require( + token.matches("[A-Za-z0-9_-]{32,128}"), + "Tokens must be distinct URL-safe strings of 32–128 characters" + ) + (token.getBytes(StandardCharsets.US_ASCII), id) + } + require(!MessageDigest.isEqual(values(0)._1, values(1)._1), "Project tokens must differ") + values + } + def project(request: Request[IO], tokens: Vector[(Array[Byte], UUID)]): Option[UUID] = { + val headers = request.headers.headers.filter(_.name.toString.equalsIgnoreCase("Authorization")) + if (headers.size != 1) None + else { + val value = headers.head.value + if (!value.startsWith("Bearer ") || value.length > 135) None + else { + val bytes = value.drop(7).getBytes(StandardCharsets.UTF_8) + tokens.find { case (expected, _) => MessageDigest.isEqual(bytes, expected) }.map(_._2) + } + } + } +} + +object Main extends IOApp.Simple { + def routes(xa: Transactor[IO], tokens: Vector[(Array[Byte], UUID)]): HttpRoutes[IO] = { + def body(request: Request[IO]): IO[Json] = request.body + .take(Domain.MaxBody.toLong + 1) + .compile + .to(Array) + .timeout(5.seconds) + .flatMap(bytes => IO(Domain.parse(bytes))) + def error(status: Int, message: String) = IO.pure( + Response[IO](Status.fromInt(status).toOption.get) + .withEntity(Json.obj("error" -> Json.fromString(message))) + ) + def protect(request: Request[IO])(operation: UUID => IO[Response[IO]]): IO[Response[IO]] = + Auth.project(request, tokens) match { + case None => error(401, "A configured project bearer token is required.") + case Some(project) => + operation(project).handleErrorWith { + case failure: Rejection => error(failure.status, failure.message) + case _: java.util.concurrent.TimeoutException => error(408, "Request timed out.") + case _: java.sql.SQLException => + error(503, "The experiment log could not complete this request. Try again shortly.") + case _ => error(500, "The request could not be completed.") + } + } + HttpRoutes.of[IO] { + case request @ GET -> Root / "project" => + protect(request)(project => Ok(Json.obj("id" -> Json.fromString(project.toString)))) + case request @ POST -> Root / "runs" => + protect(request) { project => + for { + raw <- body(request) + fields <- IO(Domain.run(raw)) + result <- Store.register(project, fields).transact(xa) + response <- if (result._2) Created(result._1.json) else Ok(result._1.json) + } yield response + } + case request @ GET -> Root / "runs" / id => + protect(request) { project => + IO(Domain.uuid(id)) + .flatMap(runId => Store.read(project, runId).transact(xa)) + .flatMap(saved => Ok(saved.json)) + } + case request @ POST -> Root / "runs" / "search" => + protect(request) { project => + for { + raw <- body(request) + fields <- IO(Domain.search(raw)) + result <- Store.search(project, fields).transact(xa) + response <- Ok(result) + } yield response + } + } + } + + def run: IO[Unit] = IO(Auth.configured).flatMap { tokens => + Database.resource.use { xa => + val api = Timeout(30.seconds)(routes(xa, tokens).orNotFound) + val check = sql"SELECT id FROM experiment_log.projects" + .query[UUID] + .to[Vector] + .transact(xa) + .flatMap(ids => + IO.raiseUnless(ids.toSet == Set(Auth.projectA, Auth.projectB))( + new IllegalStateException( + "Seeded projects do not match the configured token identities" + ) + ) + ) + check *> EmberServerBuilder + .default[IO] + .withHost(host"127.0.0.1") + .withPort(port"8080") + .withIdleTimeout(35.seconds) + .withRequestHeaderReceiveTimeout(5.seconds) + .withHttpApp(api) + .build + .useForever + } + } +} diff --git a/applications/experiment-log/src/main/scala/experimentlog/Store.scala b/applications/experiment-log/src/main/scala/experimentlog/Store.scala new file mode 100644 index 00000000..7b9f8bc5 --- /dev/null +++ b/applications/experiment-log/src/main/scala/experimentlog/Store.scala @@ -0,0 +1,126 @@ +package experimentlog + +import cats.syntax.all.* +import org.typelevel.doobie.* +import org.typelevel.doobie.implicits.* +import org.typelevel.doobie.postgres.implicits.* +import org.typelevel.doobie.postgres.circe.jsonb.implicits.* +import io.circe.Json +import java.time.Instant +import java.util.UUID + +final case class RunRow(id: UUID, requestId: UUID, title: String, config: Json, createdAt: Instant) +final case class Saved(row: RunRow, measurements: Vector[Measurement]) { + def json: Json = Json.obj( + "id" -> Json.fromString(row.id.toString), + "requestId" -> Json.fromString(row.requestId.toString), + "title" -> Json.fromString(row.title), + "config" -> row.config, + "createdAt" -> Json.fromString(row.createdAt.toString), + "measurements" -> Json.fromValues( + measurements.map(m => + Json.obj( + "name" -> Json.fromString(m.name), + "value" -> Json.fromString(Domain.decimalText(m.value)) + ) + ) + ) + ) +} + +object Store { + def header(project: UUID, id: UUID): Query0[RunRow] = + sql"""SELECT id,request_id,title,config,created_at FROM experiment_log.runs + WHERE project_id=$project AND id=$id""".query[RunRow] + + def read(project: UUID, id: UUID): ConnectionIO[Saved] = for { + row <- header(project, id).option.flatMap( + _.liftTo[ConnectionIO](Rejection("Run not found.", 404)) + ) + measurements <- sql"""SELECT name,value FROM experiment_log.measurements + WHERE project_id=$project AND run_id=$id ORDER BY name""".query[Measurement].to[Vector] + } yield Saved(row, measurements) + + def claim(project: UUID, input: RunInput, id: UUID): Query0[UUID] = { + val title = input.title + val config = input.config + val payload = input.semantic + val request = input.requestId + sql"""INSERT INTO experiment_log.runs(id,project_id,request_id,title,config,request_payload) + VALUES($id,$project,$request,$title,$config,$payload) + ON CONFLICT(project_id,request_id) DO NOTHING RETURNING id""".query[UUID] + } + + val measurementInsert = Update[(UUID, UUID, String, BigDecimal)]( + "INSERT INTO experiment_log.measurements(project_id,run_id,name,value) VALUES(?,?,?,?)" + ) + def measurements(project: UUID, id: UUID): Query0[Measurement] = + sql"""SELECT name,value FROM experiment_log.measurements + WHERE project_id=$project AND run_id=$id ORDER BY name""".query[Measurement] + + def register(project: UUID, input: RunInput): ConnectionIO[(Saved, Boolean)] = { + val id = UUID.randomUUID() + val payload = input.semantic + val request = input.requestId + for { + // One ConnectionIO, interpreted by one transact at the HTTP boundary. + claimed <- claim(project, input, id).option + result <- claimed match { + case Some(inserted) => + val rows = input.measurements.map(m => (project, inserted, m.name, m.value)) + for { + _ <- measurementInsert.updateMany(rows) + saved <- read(project, inserted) + } yield (saved, true) + case None => + // A fresh READ COMMITTED statement sees the concurrent winner after conflict wait. + for { + existing <- sql"""SELECT id,request_payload=$payload FROM experiment_log.runs + WHERE project_id=$project AND request_id=$request""".query[(UUID, Boolean)].unique + _ <- + if (existing._2) ().pure[ConnectionIO] + else + Rejection("This request UUID already records a different run.", 409) + .raiseError[ConnectionIO, Unit] + saved <- read(project, existing._1) + } yield (saved, false) + } + } yield result + } + + def searchHeaders(project: UUID, search: Search): Query0[RunRow] = { + val pattern = + "%" + search.title.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_") + "%" + val filter = search.config + val bound = search.limit.toLong + 1 + val base = fr"""SELECT id,request_id,title,config,created_at FROM experiment_log.runs + WHERE project_id=$project AND config @> $filter AND title ILIKE $pattern ESCAPE E'\\' """ + val after = + search.cursor.fold(Fragment.empty)(c => fr" AND (created_at,id) < (${c.createdAt},${c.id}) ") + (base ++ after ++ fr"ORDER BY created_at DESC,id DESC LIMIT $bound").query[RunRow] + } + + def measurementsFor(project: UUID, ids: Vector[UUID]): Query0[(UUID, String, BigDecimal)] = { + val bound = ids.toArray + sql"""SELECT run_id,name,value FROM experiment_log.measurements + WHERE project_id=$project AND run_id=ANY($bound) ORDER BY run_id,name""" + .query[(UUID, String, BigDecimal)] + } + + def search(project: UUID, input: Search): ConnectionIO[Json] = for { + rows <- searchHeaders(project, input).to[Vector] + selected = rows.take(input.limit) + children <- + if (selected.isEmpty) Vector.empty[(UUID, String, BigDecimal)].pure[ConnectionIO] + else measurementsFor(project, selected.map(_.id)).to[Vector] + grouped = children.groupMap(_._1)(row => Measurement(row._2, row._3)) + saved = selected.map(row => Saved(row, grouped.getOrElse(row.id, Vector.empty))) + next = + if (rows.size > input.limit) + selected.lastOption.map(row => Cursor(row.createdAt, row.id).encoded) + else None + } yield Json.obj( + "runs" -> Json.fromValues(saved.map(_.json)), + "nextCursor" -> next.fold(Json.Null)(Json.fromString) + ) +} diff --git a/applications/experiment-log/src/test/scala/experimentlog/DomainSuite.scala b/applications/experiment-log/src/test/scala/experimentlog/DomainSuite.scala new file mode 100644 index 00000000..81ac485c --- /dev/null +++ b/applications/experiment-log/src/test/scala/experimentlog/DomainSuite.scala @@ -0,0 +1,60 @@ +package experimentlog + +import io.circe.Json +import java.nio.charset.StandardCharsets +import java.time.Instant +import java.util.UUID + +class DomainSuite extends munit.FunSuite { + val request = "00000000-0000-4000-8000-000000000010" + def parse(value: String) = Domain.parse(value.getBytes(StandardCharsets.UTF_8)) + def payload(value: String) = + s"""{"requestId":"$request","title":"Synthetic run","config":{"optimizer":"adam","rate":0.000001},"measurements":[{"name":"accuracy","value":"$value"}]}""" + + test("exact decimal values never pass through Double or rounding") { + val saved = Domain.run(parse(payload("0.123456"))) + assertEquals(Domain.decimalText(saved.measurements.head.value), "0.123456") + for (bad <- Vector("0.1234567", "1000000.000001", "1e-6", "NaN", "01.2", "9" * 5000)) + intercept[Rejection](Domain.run(parse(payload(bad)))) + } + test("semantic equality ignores object and measurement ordering and decimal trailing zeros") { + val first = Domain.run(parse(payload("1.000000"))) + val other = Domain.run( + parse( + s"""{"measurements":[{"value":"1.0","name":"accuracy"}],"config":{"rate":0.000001,"optimizer":"adam"},"title":"Synthetic run","requestId":"$request"}""" + ) + ) + assertEquals(first.semantic, other.semantic) + assertNotEquals(first.semantic, other.copy(title = "Different").semantic) + } + test("numeric lexeme exponent and nesting bounds precede decimal expansion") { + for (bad <- Vector("1e999999999", "1e13", "9" * 100)) + intercept[Rejection](parse(s"""{"config":{"number":$bad}}""")) + intercept[Rejection](parse("[[[[[0]]]]]")) + intercept[Rejection](Domain.config(parse("""{"nested":{"key":1}}"""))) + intercept[Rejection](Domain.config(parse("""{"nothing":null}"""))) + intercept[Rejection](Domain.config(parse("""{"number":0.0000001}"""))) + } + test("duplicate keys malformed UTF-8 and unpaired surrogates fail safely") { + intercept[Rejection](parse("""{"value":1,"value":2}""")) + intercept[Rejection](Domain.parse(Array(0xc3.toByte, 0x28.toByte))) + intercept[Rejection](Domain.config(parse("""{"text":"\ud800"}"""))) + intercept[Rejection](Domain.safeText("Bad\u0085", 120)) + assertEquals(Domain.safeText("Idea 😀", 120), "Idea 😀") + } + test("measurement names are unique and caller cannot select project") { + intercept[Rejection]( + Domain.run(parse(payload("1").replace("}]}", """},{"name":"accuracy","value":"2"}]}"""))) + ) + intercept[Rejection]( + Domain.run(parse(payload("1").dropRight(1) + """, "projectId":"other"}""")) + ) + } + test("cursor is bounded canonical and preserves full timestamp and UUID") { + val cursor = Cursor(Instant.parse("2026-10-02T12:00:00.123456Z"), UUID.fromString(request)) + assertEquals(Domain.cursor(cursor.encoded), cursor) + for (bad <- Vector("!invalid", "x" * 1000, cursor.encoded + "=")) + intercept[Rejection](Domain.cursor(bad)) + intercept[Rejection](Domain.search(parse("""{"limit":21}"""))) + } +} diff --git a/applications/experiment-log/src/test/scala/experimentlog/LiveProbe.scala b/applications/experiment-log/src/test/scala/experimentlog/LiveProbe.scala new file mode 100644 index 00000000..9c942b87 --- /dev/null +++ b/applications/experiment-log/src/test/scala/experimentlog/LiveProbe.scala @@ -0,0 +1,188 @@ +package experimentlog + +import cats.effect.{IO, IOApp} +import cats.syntax.all.* +import org.typelevel.doobie.implicits.* +import org.typelevel.doobie.postgres.implicits.* +import org.typelevel.doobie.postgres.circe.jsonb.implicits.* +import io.circe.Json +import java.sql.{Connection, DriverManager} +import java.lang.reflect.{InvocationTargetException, Proxy} +import java.util.concurrent.atomic.AtomicInteger +import org.typelevel.doobie.Transactor +import java.util.{Properties, UUID} +import java.time.Instant +import javax.net.ssl.{HostnameVerifier, SSLSession} +import org.postgresql.ssl.PGjdbcHostnameVerifier + +// Test-only verifier: retain the actual TLS session/CA validation, substitute +// a wrong hostname for pgJDBC's native hostname verification control. +class WrongHostnameVerifier extends HostnameVerifier { + override def verify(hostname: String, session: SSLSession): Boolean = + PGjdbcHostnameVerifier.INSTANCE.verify("wrong-hostname.invalid", session) +} + +object LiveProbe extends IOApp.Simple { + def connect(overrides: Map[String, String]): Connection = { + val properties = new Properties() + Map( + "user" -> sys.env("PGUSER"), + "password" -> sys.env("PGPASSWORD"), + "sslmode" -> "verify-full", + "sslrootcert" -> sys.env("PGSSLROOTCERT"), + "connectTimeout" -> "10", + "socketTimeout" -> "15" + ).foreach { case (key, value) => properties.setProperty(key, value) } + overrides.foreach { case (key, value) => properties.setProperty(key, value) } + DriverManager.getConnection( + s"jdbc:postgresql://${sys.env("PGHOST")}:${sys.env.getOrElse("PGPORT", "5432")}/${sys.env.getOrElse("PGDATABASE", "postgres")}", + properties + ) + } + def driverControl(overrides: Map[String, String]): Unit = { + val connection = connect(overrides) + try { + val statement = connection.createStatement() + try { + val result = statement.executeQuery("SELECT 1") + try { require(result.next() && result.getInt(1) == 1) } + finally result.close() + } finally statement.close() + } finally connection.close() + } + def messages(error: Throwable): String = Iterator + .iterate(Option(error))(_.flatMap(e => Option(e.getCause))) + .takeWhile(_.nonEmpty) + .flatMap(_.map(_.getMessage)) + .mkString(" | ") + + def run: IO[Unit] = for { + _ <- IO.blocking(driverControl(Map.empty)) + _ <- IO.println("pgJDBC verify-full positive connection passed") + wrongCa <- IO + .blocking(driverControl(Map("sslrootcert" -> "/etc/ssl/certs/ca-certificates.crt"))) + .attempt + _ <- IO { + require(wrongCa.isLeft, "wrong CA must fail"); + require( + messages(wrongCa.swap.toOption.get).toLowerCase + .matches("(?s).*(certificate|pkix|trust anchor|certification path).*"), + "CA control must fail for certificate validation" + ) + } + _ <- IO.println("pgJDBC wrong CA reached a certificate-specific failure") + wrongName <- IO + .blocking(driverControl(Map("sslhostnameverifier" -> classOf[WrongHostnameVerifier].getName))) + .attempt + _ <- IO { + require(wrongName.isLeft, "wrong hostname must fail"); + require( + messages(wrongName.swap.toOption.get).toLowerCase.matches("(?s).*(hostname|host name).*"), + "hostname control must fail during certificate name validation" + ) + } + _ <- IO.println( + "pgJDBC actual TLS session + native verifier substituted-hostname negative passed; production verifier unchanged" + ) + _ <- IO + .blocking(connect(Map.empty)) + .bracket { connection => + val count = new AtomicInteger(0) + val wrapped = Proxy + .newProxyInstance( + classOf[Connection].getClassLoader, + Array(classOf[Connection]), + (_, method, arguments) => { + if (method.getName == "prepareStatement") count.incrementAndGet() + try method.invoke(connection, Option(arguments).getOrElse(Array.empty[AnyRef])*) + catch { case error: InvocationTargetException => throw error.getCause } + } + ) + .asInstanceOf[Connection] + val xa = Transactor.fromConnection[IO](wrapped, logHandler = None) + Store + .search(Auth.projectA, Search(Json.obj(), "", 20, None)) + .transact(xa) + .flatMap(result => + IO { + require( + result.hcursor.downField("runs").focus.flatMap(_.asArray).exists(_.nonEmpty), + "count control requires persisted fixture rows" + ) + require( + count.get == 2, + s"Expected exactly two prepared search statements, got ${count.get}" + ) + println( + "Observed exactly 2 JDBC prepared statements for populated limit-20 search; no SQL parameter logging" + ) + } + ) + }(connection => IO.blocking(connection.close())) + _ <- Database.resource.use { xa => + val project = Auth.projectA + val id = UUID.randomUUID() + val input = RunInput( + UUID.randomUUID(), + "Typecheck fixture", + Json.obj("optimizer" -> Json.fromString("adam")), + Vector(Measurement("accuracy", BigDecimal("0.123456"))) + ) + val checks = Vector( + "read header" -> Store.header(project, id).analysis, + "measurement decimal" -> Store.measurements(project, id).analysis, + "bounded measurement batch" -> Store.measurementsFor(project, Vector(id)).analysis, + "retained insert returning" -> Store.claim(project, input, id).analysis, + "child insert exact numeric" -> Store.measurementInsert + .toUpdate0((project, id, "accuracy", BigDecimal("0.123456"))) + .analysis, + "containment search" -> Store + .searchHeaders(project, Search(input.config, "%_", 5, None)) + .analysis, + "cursor search" -> Store + .searchHeaders( + project, + Search( + input.config, + "", + 5, + Some(Cursor(Instant.parse("2026-10-02T12:00:00.123456Z"), id)) + ) + ) + .analysis + ) + for { + _ <- checks.traverse_ { case (name, check) => + check + .transact(xa) + .flatMap(a => + IO { + require( + a.alignmentErrors.isEmpty, + s"$name SQL type mismatch: ${a.alignmentErrors}" + ); + println(s"doobie live analysis passed: $name") + } + ) + } + version <- sql"SELECT version()".query[String].unique.transact(xa) + _ <- IO.println(version) + definition <- + sql"SELECT indexdef FROM pg_indexes WHERE schemaname='experiment_log' AND indexname='runs_config_gin'" + .query[String] + .unique + .transact(xa) + _ <- IO.println(definition) + filter = input.config + plan <- + sql"EXPLAIN SELECT id FROM experiment_log.runs WHERE project_id=$project AND config @> $filter" + .query[String] + .to[Vector] + .transact(xa) + _ <- IO.println( + "Actual small-fixture EXPLAIN (planner not forced):\n" + plan.mkString("\n") + ) + } yield () + } + } yield () +} diff --git a/applications/experiment-log/tests/acceptance.py b/applications/experiment-log/tests/acceptance.py new file mode 100644 index 00000000..2296303c --- /dev/null +++ b/applications/experiment-log/tests/acceptance.py @@ -0,0 +1,372 @@ +"""Opt-in real Cloud + running HTTP API acceptance; never run in CI.""" + +import copy +import json +import os +import time +import unittest +import uuid +from concurrent.futures import ThreadPoolExecutor +from decimal import Decimal + +import httpx +import psycopg +from psycopg.types.json import Jsonb + +BASE = os.environ.get("BASE_URL", "http://127.0.0.1:8080") +A = "00000000-0000-4000-8000-000000000001" +B = "00000000-0000-4000-8000-000000000002" + + +def connection(owner=False): + return psycopg.connect( + host=os.environ["PGHOST"], + port=os.environ.get("PGPORT", "5432"), + dbname=os.environ.get("PGDATABASE", "postgres"), + user=os.environ["TEST_OWNER_USER"] if owner else os.environ["PGUSER"], + password=os.environ["TEST_OWNER_PASSWORD"] + if owner + else os.environ["PGPASSWORD"], + sslmode="verify-full", + sslrootcert=os.environ["PGSSLROOTCERT"], + connect_timeout=10, + options="-c statement_timeout=10000 -c lock_timeout=8000", + ) + + +def request(method, path, data=None, project="A", raw=None): + with httpx.Client(base_url=BASE, timeout=35) as client: + headers = {"Authorization": "Bearer " + os.environ[f"PROJECT_{project}_TOKEN"]} + if raw is not None: + headers["Content-Type"] = "application/json" + return client.request(method, path, content=raw, headers=headers) + return client.request(method, path, json=data, headers=headers) + + +def payload(title="Synthetic optimizer comparison", config=None): + return { + "requestId": str(uuid.uuid4()), + "title": title, + "config": {"optimizer": "adam", "enabled": True, "rate": 0.000001} + if config is None + else config, + "measurements": [ + {"name": "accuracy", "value": "0.123456"}, + {"name": "loss", "value": "-0.000001"}, + ], + } + + +class CloudAcceptance(unittest.TestCase): + def create(self, body=None, project="A"): + body = payload() if body is None else body + response = request("POST", "/runs", body, project) + self.assertEqual(response.status_code, 201, response.text) + return response.json() + + def test_01_exact_decimal_and_project_scope(self): + body = payload() + body["measurements"] += [ + {"name": "minimum", "value": "-1000000"}, + {"name": "maximum", "value": "1000000"}, + ] + saved = self.create(body) + values = {m["name"]: m["value"] for m in saved["measurements"]} + self.assertEqual( + values, + { + "accuracy": "0.123456", + "loss": "-0.000001", + "minimum": "-1000000", + "maximum": "1000000", + }, + ) + self.assertEqual( + json.loads(json.dumps(saved), parse_float=Decimal)["config"]["rate"], + Decimal("0.000001"), + ) + self.assertEqual(request("GET", "/runs/" + saved["id"]).json(), saved) + self.assertEqual( + request("GET", "/runs/" + saved["id"], project="B").status_code, 404 + ) + other = self.create(body, "B") + self.assertNotEqual(other["id"], saved["id"]) + self.assertEqual(request("GET", "/project", project="B").json(), {"id": B}) + self.assertEqual(httpx.get(BASE + "/project").status_code, 401) + self.assertEqual( + httpx.get( + BASE + "/project", headers={"Authorization": "Bearer wrong"} + ).status_code, + 401, + ) + self.assertEqual(request("GET", "/runs/not-a-uuid").status_code, 422) + + def test_02_semantic_replay_and_conflict(self): + body = payload(config={"b": True, "a": 1}) + body["measurements"][0]["value"] = "1.000000" + saved = self.create(body) + reordered = copy.deepcopy(body) + reordered["config"] = {"a": 1.0, "b": True} + reordered["measurements"].reverse() + next(m for m in reordered["measurements"] if m["name"] == "accuracy")[ + "value" + ] = "1.0" + response = request("POST", "/runs", reordered) + self.assertEqual(response.status_code, 200, response.text) + self.assertEqual(response.json(), saved) + reordered["title"] += " changed" + self.assertEqual(request("POST", "/runs", reordered).status_code, 409) + + def test_03_http_boundaries(self): + for change in [ + {"projectId": A}, + {"config": {"x": None}}, + {"config": {"x": []}}, + {"config": {"x": {}}}, + {"config": {"rate": 0.0000001}}, + {"title": "\u0085bad"}, + {"title": "\ud800"}, + {"measurements": [{"name": "metric", "value": "0.1234567"}]}, + {"measurements": [{"name": "metric", "value": "1e3"}]}, + {"measurements": [{"name": "metric", "value": 1}]}, + {"measurements": []}, + ]: + body = payload() + body.update(change) + self.assertEqual( + request("POST", "/runs", raw=json.dumps(body).encode()).status_code, + 422, + str(change), + ) + malformed = [b'{"x":1,"x":2}', b"\xff"] + for raw in malformed: + self.assertEqual(request("POST", "/runs", raw=raw).status_code, 400) + out_of_bounds = [b'{"x":1e999999}', b'{"x":[[[[[1]]]]]}', b" " * 16385] + for raw in out_of_bounds: + self.assertEqual(request("POST", "/runs", raw=raw).status_code, 422) + body = payload(config={"x": 0.000001}) + self.create(body) + + def competing(self, conflicting): + first = payload() + second = copy.deepcopy(first) + if conflicting: + second["title"] += " other" + with connection(True) as owner: + owner.execute("LOCK TABLE experiment_log.runs IN SHARE ROW EXCLUSIVE MODE") + with ThreadPoolExecutor(max_workers=2) as pool: + calls = [ + pool.submit(request, "POST", "/runs", body) + for body in [first, second] + ] + deadline = time.monotonic() + 5 + blocked = 0 + while time.monotonic() < deadline: + owner.execute("SELECT pg_stat_clear_snapshot()") + blocked = owner.execute( + "SELECT count(*) FROM pg_stat_activity WHERE usename='experiment_app' AND cardinality(pg_blocking_pids(pid)) > 0" + ).fetchone()[0] + if blocked >= 2: + break + time.sleep(0.05) + self.assertGreaterEqual( + blocked, + 2, + "both independent HTTP transactions must contend on a real DB lock", + ) + owner.commit() + responses = [call.result() for call in calls] + self.assertEqual( + sorted(r.status_code for r in responses), + [201, 409] if conflicting else [200, 201], + ) + success = next(r.json() for r in responses if r.status_code == 201) + if not conflicting: + self.assertEqual(responses[0].json(), responses[1].json()) + with connection(True) as owner: + self.assertEqual( + owner.execute( + "SELECT count(*) FROM experiment_log.runs WHERE project_id=%s AND request_id=%s", + (A, first["requestId"]), + ).fetchone()[0], + 1, + ) + self.assertEqual( + owner.execute( + "SELECT count(*) FROM experiment_log.measurements WHERE run_id=%s", + (success["id"],), + ).fetchone()[0], + 2, + ) + + def test_04_competing_matching_retained_request(self): + self.competing(False) + + def test_05_competing_conflicting_retained_request(self): + self.competing(True) + + def test_06_after_parent_and_first_child_rollback(self): + body = payload() + body["measurements"] = [ + {"name": "a_first", "value": "1"}, + {"name": "z_forced_failure", "value": "2"}, + ] + with connection(True) as owner: + owner.execute("""CREATE FUNCTION experiment_log.acceptance_fault() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN + IF NEW.name='z_forced_failure' AND EXISTS (SELECT 1 FROM experiment_log.measurements WHERE run_id=NEW.run_id AND name='a_first') THEN + RAISE EXCEPTION 'Controlled child failure' USING ERRCODE='23514'; END IF; RETURN NEW; END $$""") + owner.execute( + "CREATE TRIGGER acceptance_fault AFTER INSERT ON experiment_log.measurements FOR EACH ROW EXECUTE FUNCTION experiment_log.acceptance_fault()" + ) + try: + response = request("POST", "/runs", body) + self.assertEqual(response.status_code, 503, response.text) + with connection(True) as owner: + self.assertEqual( + owner.execute( + "SELECT count(*) FROM experiment_log.runs WHERE request_id=%s", + (body["requestId"],), + ).fetchone()[0], + 0, + ) + self.assertEqual( + owner.execute( + "SELECT count(*) FROM experiment_log.measurements m JOIN experiment_log.runs r ON r.id=m.run_id AND r.project_id=m.project_id WHERE r.request_id=%s", + (body["requestId"],), + ).fetchone()[0], + 0, + ) + finally: + with connection(True) as owner: + owner.execute( + "DROP TRIGGER acceptance_fault ON experiment_log.measurements" + ) + owner.execute("DROP FUNCTION experiment_log.acceptance_fault()") + self.create(body) + + def test_07_immutable_role_and_initial_transaction(self): + saved = self.create() + denied = [ + "UPDATE experiment_log.runs SET title='changed'", + "DELETE FROM experiment_log.runs", + "CREATE TABLE experiment_log.forbidden(id int)", + "CREATE TEMP TABLE forbidden(id int)", + ] + for statement in denied: + with connection() as runtime: + with self.assertRaises(psycopg.errors.InsufficientPrivilege): + runtime.execute(statement) + with connection() as runtime: + with self.assertRaises(psycopg.errors.CheckViolation): + runtime.execute( + "INSERT INTO experiment_log.measurements(project_id,run_id,name,value) VALUES(%s,%s,'late',1)", + (A, saved["id"]), + ) + self.assertEqual(request("GET", "/runs/" + saved["id"]).json(), saved) + + def test_08_database_complete_set_and_no_numeric_rounding(self): + body = payload() + semantic = { + "title": body["title"], + "config": body["config"], + "measurements": {"accuracy": "0.123456", "loss": "-0.000001"}, + } + with connection() as runtime: + with self.assertRaises(psycopg.errors.CheckViolation): + runtime.execute( + "INSERT INTO experiment_log.runs(id,project_id,request_id,title,config,request_payload) VALUES(%s,%s,%s,%s,%s,%s)", + ( + uuid.uuid4(), + A, + body["requestId"], + body["title"], + Jsonb(body["config"]), + Jsonb(semantic), + ), + ) + runtime.commit() + with connection(True) as owner: + self.assertEqual( + owner.execute( + "SELECT count(*) FROM experiment_log.runs WHERE request_id=%s", + (body["requestId"],), + ).fetchone()[0], + 0, + ) + with connection(True) as owner: + run_id = uuid.uuid4() + owner.execute( + "INSERT INTO experiment_log.runs(id,project_id,request_id,title,config,request_payload) VALUES(%s,%s,%s,%s,%s,%s)", + ( + run_id, + A, + uuid.uuid4(), + body["title"], + Jsonb(body["config"]), + Jsonb(semantic), + ), + ) + with self.assertRaises(psycopg.errors.CheckViolation): + owner.execute( + "INSERT INTO experiment_log.measurements(project_id,run_id,name,value) VALUES(%s,%s,'accuracy',0.1234567)", + (A, run_id), + ) + + def test_09_jsonb_literal_search_and_keyset(self): + marker = uuid.uuid4().hex + created = [ + self.create( + payload( + f"{marker} literal%_ item {i}", + {"suite": marker, "optimizer": "adam", "enabled": True}, + ) + ) + for i in range(5) + ] + first = request( + "POST", + "/runs/search", + {"config": {"suite": marker, "enabled": True}, "title": "%_", "limit": 2}, + ) + self.assertEqual(first.status_code, 200, first.text) + seen = [] + cursor = None + for _ in range(4): + fields = {"config": {"suite": marker}, "limit": 2} + if cursor: + fields["cursor"] = cursor + response = request("POST", "/runs/search", fields) + self.assertEqual(response.status_code, 200, response.text) + data = response.json() + seen.extend(row["id"] for row in data["runs"]) + cursor = data["nextCursor"] + if cursor is None: + break + self.assertEqual(len(seen), 5) + self.assertEqual(set(seen), {row["id"] for row in created}) + self.assertEqual( + request( + "POST", "/runs/search", {"config": {"suite": marker, "enabled": False}} + ).json()["runs"], + [], + ) + empty_filter = request( + "POST", "/runs/search", {"config": {}, "limit": 20} + ).json() + self.assertGreater(len(empty_filter["runs"]), 0) + self.assertLessEqual(len(empty_filter["runs"]), 20) + for fields in [ + {"limit": 21}, + {"cursor": "invalid"}, + {"config": []}, + {"config": {"x": None}}, + ]: + self.assertEqual(request("POST", "/runs/search", fields).status_code, 422) + other = request( + "POST", "/runs/search", {"config": {"suite": marker}}, project="B" + ) + self.assertEqual(other.json()["runs"], []) + + +if __name__ == "__main__": + unittest.main(verbosity=2) diff --git a/applications/experiment-log/tests/requirements.in b/applications/experiment-log/tests/requirements.in new file mode 100644 index 00000000..f6139e6a --- /dev/null +++ b/applications/experiment-log/tests/requirements.in @@ -0,0 +1,3 @@ +httpx==0.28.1 +psycopg[binary]==3.3.6 +ruff==0.16.10 diff --git a/applications/experiment-log/tests/requirements.txt b/applications/experiment-log/tests/requirements.txt new file mode 100644 index 00000000..63ec0f80 --- /dev/null +++ b/applications/experiment-log/tests/requirements.txt @@ -0,0 +1,10 @@ +anyio==4.15.1 +certifi==2026.7.22 +h11==0.16.0 +httpcore==1.0.9 +httpx==0.28.1 +idna==3.20 +psycopg==3.3.6 +psycopg-binary==3.3.6 +ruff==0.16.10 +typing_extensions==4.16.0 From 2ae987f340186383e6acafad8138b0dc0b969181 Mon Sep 17 00:00:00 2001 From: sdairs Date: Fri, 2 Oct 2026 21:31:28 +0100 Subject: [PATCH 2/2] docs: remove public beta qualifier --- applications/experiment-log/README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/applications/experiment-log/README.md b/applications/experiment-log/README.md index 5f639ffb..8a90b6e6 100644 --- a/applications/experiment-log/README.md +++ b/applications/experiment-log/README.md @@ -1,6 +1,6 @@ # Experiment log -A small Scala API records supplied experiment results in ClickHouse Managed Postgres (public beta). Two seeded projects have separate bearer tokens. Each immutable run contains a title, a bounded flat JSONB configuration and 1–10 exact decimal measurements. These are synthetic examples, not model training or scientific validation. +A small Scala API records supplied experiment results in ClickHouse Managed Postgres. Two seeded projects have separate bearer tokens. Each immutable run contains a title, a bounded flat JSONB configuration and 1–10 exact decimal measurements. These are synthetic examples, not model training or scientific validation. ## Requirements and verified versions