diff --git a/.github/workflows/task-dependencies.yml b/.github/workflows/task-dependencies.yml new file mode 100644 index 00000000..3efb99e5 --- /dev/null +++ b/.github/workflows/task-dependencies.yml @@ -0,0 +1,26 @@ +name: Task dependencies +on: + push: + paths: ['applications/task-dependencies/**', '.github/workflows/task-dependencies.yml'] + pull_request: + paths: ['applications/task-dependencies/**', '.github/workflows/task-dependencies.yml'] +permissions: + contents: read +jobs: + checks: + runs-on: ubuntu-24.04 + defaults: + run: + working-directory: applications/task-dependencies + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-go@v5 + with: + go-version: '1.27.1' + cache-dependency-path: applications/task-dependencies/go.sum + - run: go mod download + - run: go mod verify + - run: test -z "$(gofmt -l cmd internal)" + - run: go test -tags nomsgpack ./... + - run: go vet -tags nomsgpack ./... + - run: go build -tags nomsgpack ./cmd/server diff --git a/applications/task-dependencies/.env.example b/applications/task-dependencies/.env.example new file mode 100644 index 00000000..eaf7cf00 --- /dev/null +++ b/applications/task-dependencies/.env.example @@ -0,0 +1,9 @@ +PGHOST=your-service-hostname +PGPORT=5432 +PGDATABASE=postgres +PGUSER=tasks_runtime +PGPASSWORD=generate-a-private-runtime-password +PGSSLMODE=verify-full +PGSSLROOTCERT=/absolute/path/official-ca.pem +TASKS_NORTH_TOKEN=generate-a-distinct-private-32-byte-token +TASKS_SOUTH_TOKEN=generate-another-private-32-byte-token diff --git a/applications/task-dependencies/.gitignore b/applications/task-dependencies/.gitignore new file mode 100644 index 00000000..f01363bd --- /dev/null +++ b/applications/task-dependencies/.gitignore @@ -0,0 +1,6 @@ +/bin/ +/.local/ +/.venv/ +__pycache__/ +.env +*.pem diff --git a/applications/task-dependencies/README.md b/applications/task-dependencies/README.md new file mode 100644 index 00000000..147ee42d --- /dev/null +++ b/applications/task-dependencies/README.md @@ -0,0 +1,158 @@ +# Task dependencies with Go, Gin and GORM + +A bounded account-scoped task API on **ClickHouse Managed Postgres**. Create tasks inside a seeded project, add/remove prerequisites, and complete a task once its current prerequisites are done. The project lock makes one graph authoritative while competing changes are checked and committed. + +Pinned stack: Go 1.27.1, Gin 1.12.0, GORM 1.31.2, GORM Postgres driver 1.6.3 and pgx 5.11.0. Native Linux ARM64 was tested; PostgreSQL 18 is the Cloud major. No browser interface, project creation, task deletion or reopening route is included. This is a loopback API with two synthetic accounts, not an identity provider. + +## Transaction and graph rules + +Every graph write begins a GORM transaction and locks the account-scoped project `FOR UPDATE`. It reads at most 100 tasks and 300 edges, validates the requested change, then updates the project revision and writes through **that same `tx` handle**. A task points toward its prerequisites; adding `A → B` cycles if B already reaches A. Opposing concurrent edges cannot both validate against an older graph. + +Completion checks current prerequisites under the same lock as edge changes. A completed task cannot gain a new prerequisite. Existing edge additions return the original edge, removal is repeat-safe, and repeated completion returns the saved task without changing its timestamp/revision. Removing an existing prerequisite from a completed task is allowed; it does not reopen the task. Task creation generates a new UUID each time and is **not** repeat-safe after an ambiguous response: read the graph before deciding whether another task is needed. + +The database enforces unique/no-self edges, composite foreign keys for both endpoints within the same project, and 100-task/300-edge insertion bounds. A project revision stays between 0 and 1,000,000,000; further changes conflict at the cap. Revisions are server output, not client optimistic-edit tokens. Snapshot reads hold a project `FOR SHARE` lock through their ordered task/edge queries, keeping the header and graph consistent with writers. Project lists use a UUID cursor, stable ascending order and a fixed 20-row limit. Graph reads use hard limits and reject unexpected oversized stored graphs. + +No runtime `AutoMigrate`, `Save` insert fallback, or zero-value struct update is used. Explicit map/column updates check `Error` and `RowsAffected`; completion sets `done=true` and its timestamp with `RETURNING`. + +## Build in native Linux + +Use Go 1.27.1, Python 3.12+, `psql` and `clickhousectl` in a Linux filesystem. The `nomsgpack` build tag omits unused Gin MessagePack support; JSON is the only accepted write format. + +```bash +go mod download +go mod verify +go test -tags nomsgpack ./... +go vet -tags nomsgpack ./... +mkdir -p bin +go build -tags nomsgpack -o bin/task-api ./cmd/server +``` + +`go.mod` and `go.sum` pin dependencies. Database-free units cover graph reachability, Unicode/UUID boundaries, strict JSON, credential configuration and foreign origins. Real Cloud controls skip unless explicitly enabled. + +## Create a disposable Cloud fixture + +Authenticate the CLI privately with your own API key. Choose your organization explicitly and confirm current region/size availability and pricing in the [Managed Postgres documentation](https://clickhouse.com/docs/products/managed-postgres/). This c6gd.large AWS fixture has no HA and incurs charges; do not use an existing production database for destructive failure controls. + +```bash +export ORG_ID=your-organization-id +umask 077 +mkdir -p .local/private +clickhousectl cloud postgres create \ + --name task-dependencies-demo --provider aws --region us-east-1 \ + --size c6gd.large --pg-version 18 --ha-type none \ + --org-id "$ORG_ID" --json > .local/private/create.json +export PG_ID=$(python3 -c 'import json; print(json.load(open(".local/private/create.json"))["id"])') +clickhousectl cloud postgres get "$PG_ID" --org-id "$ORG_ID" --json +# Repeat get until state is running before downloading its CA or connecting. +clickhousectl cloud postgres certs get "$PG_ID" --org-id "$ORG_ID" \ + --output .local/private/ca.pem +``` + +The creation receipt contains administrator credentials once; `get` does not recover the password. Keep four environment files outside Git, mode 600: + +| File | Contents | +| --- | --- | +| `admin.env` | Common fields, administrator `PGUSER`/`PGPASSWORD`, generated `TASKS_OWNER_PASSWORD` and `TASKS_RUNTIME_PASSWORD` | +| `migration.env` | Common fields, `PGUSER=tasks_owner`, its password | +| `runtime.env` | Common fields, `PGUSER=tasks_runtime`, its password, distinct `TASKS_NORTH_TOKEN` and `TASKS_SOUTH_TOKEN` | +| `test.env` | `TEST_OWNER_USER=tasks_owner`, `TEST_OWNER_PASSWORD` only | + +Common fields: `PGHOST`, `PGPORT=5432`, `PGDATABASE=postgres`, `PGSSLMODE=verify-full`, `PGSSLROOTCERT=/absolute/path/official-ca.pem`, `PGCONNECT_TIMEOUT=10`. Generate private passwords/tokens using `python3 -c 'import secrets; print(secrets.token_urlsafe(32))'`. Bearer tokens must differ and contain 32–256 bytes; they map to fixed synthetic north/south account UUIDs on the server. Never put credentials in URLs or committed source. + +In a setup terminal: + +```bash +set -a; source .local/private/admin.env; set +a +psql -X -f sql/bootstrap.sql +set -a; source .local/private/migration.env; set +a +bash scripts/migrate.sh up +psql -X -f sql/grants.sql +psql -X -f sql/seed.sql +``` + +Migration `up` is advisory-locked and checks the recorded SHA256. Repeating it verifies the checksum. Repeating seed inserts no duplicate projects/tasks/edges. `down` destroys this app’s data/objects while preserving the administrator-created schema and clears its migration version; follow with up/grants/seed only on a disposable fixture. Do not edit applied migrations. + +## Run and use the API + +Open a new terminal that has not sourced setup credentials. Export **only runtime credentials**: + +```bash +set -a; source .local/private/runtime.env; set +a +./bin/task-api +``` + +The actual GORM driver receives a pgx/database/sql connection configured with the official CA and mandatory `verify-full`. Certificate and hostname checks stay enabled. The pool has four connections, 30-second idle and 15-minute connection lifetimes; connection timeout is five seconds, statements ten seconds and lock waits five seconds. HTTP time/header limits and an 8 KiB JSON write-body limit bound requests. Graceful shutdown joins before closing the pool. + +From a private client terminal with the north token in `TOKEN`: + +```bash +export TOKEN="$TASKS_NORTH_TOKEN" +curl -s http://127.0.0.1:8090/api/projects -H "Authorization: Bearer $TOKEN" +curl -s http://127.0.0.1:8090/api/projects/20000000-0000-4000-8000-000000000001 \ + -H "Authorization: Bearer $TOKEN" +curl -s http://127.0.0.1:8090/api/projects/20000000-0000-4000-8000-000000000001/tasks \ + -H "Authorization: Bearer $TOKEN" -H 'Content-Type: application/json' \ + --data-binary '{"title":"Check launch checklist"}' +``` + +Use returned task UUIDs in these routes; writes require `application/json` with exactly the documented object fields, no duplicates: + +| Method/path below `/api/projects/:project` | Body/result | +| --- | --- | +| `GET` | Coherent project, ordered tasks and edges | +| `POST /tasks` | `{"title":"Review"}`; new task, HTTP 201 | +| `POST /tasks/:task/prerequisites` | `{"prerequisite_id":"UUID"}`; existing/new edge, HTTP 200 | +| `DELETE /tasks/:task/prerequisites/:prerequisite` | Repeat-safe HTTP 204 | +| `POST /tasks/:task/complete` | `{}`; completed/existing task, HTTP 200 | + +Cycles, unresolved prerequisites, completed-task additions and caps return 409 with a public code. Cross-account or missing project/task UUIDs return 404; malformed paths/JSON return 400 and invalid titles 422. UUIDs are canonicalized. Titles reject Unicode controls and replacement characters, conservatively rejecting JSON lone surrogates decoded as U+FFFD; valid paired non-BMP text works. Database/internal errors return a generic 503 without raw SQL or credentials. + +## Trust and omissions + +All routes derive account identity from server-mapped bearer credentials; submitted IDs cannot choose an account. There are no ambient cookies or CORS permissions. Unsafe requests with a foreign `Origin` are rejected. Tokens have no expiry/revocation store: rotate the server’s token configuration and restart to revoke them. Production requires HTTPS, real credential issuance/rotation and appropriate perimeter limits; keep this demonstration on loopback. + +The shared runtime role is trusted across both accounts. It can issue permitted direct SQL across accounts and bypass API cycle/prerequisite/monotonic-completion checks. There is no RLS, universal DAG constraint or database-enforced account authorization. Database endpoint/bound constraints still apply. The separate migration owner can change schema and data; runtime cannot assume that role, create schema objects/temp tables, change task titles/account ownership or delete tasks/projects. + +## Actual Cloud acceptance and restart + +Use only your disposable fixture. In a test terminal, export runtime plus test credentials, never into the running API: + +```bash +python3 -m venv .venv +.venv/bin/pip install -r tests/requirements.txt +set -a; source .local/private/runtime.env; source .local/private/test.env; set +a +.venv/bin/python tests/acceptance.py +``` + +Nine HTTP cases exercise workflow/retries, direct/long cycles, observed opposing-edge contention, completion versus a new prerequisite, rollback after the revision UPDATE, account scope/composite keys, graph quotas, input/revision bounds and runtime privileges. Helpers create/delete their own owner-controlled fixtures and temporary failure trigger. They are not routes. + +For actual GORM TLS and a paused SHARE reader versus an independent waiting HTTP writer, create the dedicated empty coherence fixture once using the test owner: + +```bash +export PROBE_PROJECT_ID=20000000-0000-4000-8000-000000000099 +PGUSER="$TEST_OWNER_USER" PGPASSWORD="$TEST_OWNER_PASSWORD" psql -X -c \ + "INSERT INTO task_dependencies.projects(id,account_id,name) VALUES ('$PROBE_PROJECT_ID','10000000-0000-4000-8000-000000000001','Coherence fixture')" +TASKS_CLOUD_TESTS=1 go test -tags nomsgpack ./internal/board -run TestCloud -count=1 -v +``` + +The TLS controls use the same GORM driver construction: an empty trust store and a substituted TLS server name must fail specifically for certificate/hostname verification. Application configuration uses the unchanged official CA/hostname. The coherence control observes a writer blocked by the exact reader backend, then compares old and subsequent complete snapshots. Repeating that test requires owner cleanup of its task and resetting the dedicated project revision. + +After the tests, replace the exact app process and compare durable API state: + +```bash +API_PID=your-api-pid EVIDENCE_DIR=/tmp/task-dependencies-evidence \ + .venv/bin/python tests/restart.py +``` + +This one-time helper creates its own persistence project, completes three tasks and two prerequisite edges, verifies executable/cwd, waits for the original process to exit, and starts a runtime-only replacement. It compares the exact project/revision, task completion timestamps and edges plus repeat completion. It records `persistence.json` and `replacement.pid`; stop the replacement during cleanup. Local CI needs no Cloud credentials. + +## Cleanup + +Stop the API process, then delete only the dedicated Cloud service you created and verify that exact ID is absent: + +```bash +clickhousectl cloud postgres delete "$PG_ID" --org-id "$ORG_ID" --json +clickhousectl cloud postgres list --org-id "$ORG_ID" --json +``` + +Delete promptly after review to end fixture billing. If retaining Cloud while removing app objects, use an administrator terminal that explicitly restores receipt `PGUSER`/`PGPASSWORD` before dropping this example’s schema/roles. Do not clean up an unrelated or existing service. diff --git a/applications/task-dependencies/cmd/server/main.go b/applications/task-dependencies/cmd/server/main.go new file mode 100644 index 00000000..04ea5c45 --- /dev/null +++ b/applications/task-dependencies/cmd/server/main.go @@ -0,0 +1,52 @@ +package main + +import ( + "context" + "errors" + "github.com/ClickHouse/examples/applications/task-dependencies/internal/board" + "github.com/ClickHouse/examples/applications/task-dependencies/internal/httpapi" + "log" + "net/http" + "os" + "os/signal" + "syscall" + "time" +) + +func main() { + credentials, err := httpapi.Credentials(os.Getenv("TASKS_NORTH_TOKEN"), os.Getenv("TASKS_SOUTH_TOKEN")) + if err != nil { + log.Fatal(err) + } + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + db, pool, err := board.Open(ctx) + cancel() + if err != nil { + log.Fatal("Verified database connection failed; check private setup.") + } + defer pool.Close() + var version int + if err = db.Raw("SELECT version FROM task_dependencies.schema_migrations WHERE version = 1").Scan(&version).Error; err != nil || version != 1 { + log.Fatal("Apply versioned migrations before startup.") + } + server := &http.Server{Addr: "127.0.0.1:8090", Handler: httpapi.Router(board.Store{DB: db}, credentials), ReadHeaderTimeout: 5 * time.Second, ReadTimeout: 10 * time.Second, WriteTimeout: 15 * time.Second, IdleTimeout: 60 * time.Second, MaxHeaderBytes: 16384} + stop, stopCancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) + defer stopCancel() + joined := make(chan struct{}) + go func() { + defer close(joined) + <-stop.Done() + shutdown, done := context.WithTimeout(context.Background(), 15*time.Second) + defer done() + if err := server.Shutdown(shutdown); err != nil { + server.Close() + } + }() + log.Print("Task dependency API listening on loopback port 8090.") + if err = server.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { + stopCancel() + <-joined + log.Fatal("Local HTTP listener failed.") + } + <-joined +} diff --git a/applications/task-dependencies/go.mod b/applications/task-dependencies/go.mod new file mode 100644 index 00000000..e7c42afa --- /dev/null +++ b/applications/task-dependencies/go.mod @@ -0,0 +1,48 @@ +module github.com/ClickHouse/examples/applications/task-dependencies + +go 1.27.1 + +require ( + github.com/gin-gonic/gin v1.12.0 + github.com/jackc/pgx/v5 v5.11.0 + gorm.io/driver/postgres v1.6.3 + gorm.io/gorm v1.31.2 +) + +require ( + github.com/bytedance/gopkg v0.1.3 // indirect + github.com/bytedance/sonic v1.15.0 // indirect + github.com/bytedance/sonic/loader v0.5.0 // indirect + github.com/cloudwego/base64x v0.1.6 // indirect + github.com/gabriel-vasile/mimetype v1.4.12 // indirect + github.com/gin-contrib/sse v1.1.0 // indirect + github.com/go-playground/locales v0.14.1 // indirect + github.com/go-playground/universal-translator v0.18.1 // indirect + github.com/go-playground/validator/v10 v10.30.1 // indirect + github.com/goccy/go-json v0.10.5 // indirect + github.com/goccy/go-yaml v1.19.2 // indirect + github.com/jackc/pgpassfile v1.0.0 // indirect + github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect + github.com/jackc/puddle/v2 v2.2.2 // indirect + github.com/jinzhu/inflection v1.0.0 // indirect + github.com/jinzhu/now v1.1.5 // indirect + github.com/json-iterator/go v1.1.12 // indirect + github.com/klauspost/cpuid/v2 v2.3.0 // indirect + github.com/leodido/go-urn v1.4.0 // indirect + github.com/mattn/go-isatty v0.0.20 // indirect + github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect + github.com/modern-go/reflect2 v1.0.2 // indirect + github.com/pelletier/go-toml/v2 v2.2.4 // indirect + github.com/quic-go/qpack v0.6.0 // indirect + github.com/quic-go/quic-go v0.59.0 // indirect + github.com/twitchyliquid64/golang-asm v0.15.1 // indirect + github.com/ugorji/go/codec v1.3.1 // indirect + go.mongodb.org/mongo-driver/v2 v2.5.0 // indirect + golang.org/x/arch v0.22.0 // indirect + golang.org/x/crypto v0.48.0 // indirect + golang.org/x/net v0.51.0 // indirect + golang.org/x/sync v0.19.0 // indirect + golang.org/x/sys v0.41.0 // indirect + golang.org/x/text v0.34.0 // indirect + google.golang.org/protobuf v1.36.10 // indirect +) diff --git a/applications/task-dependencies/go.sum b/applications/task-dependencies/go.sum new file mode 100644 index 00000000..1ce35ad2 --- /dev/null +++ b/applications/task-dependencies/go.sum @@ -0,0 +1,112 @@ +github.com/bytedance/gopkg v0.1.3 h1:TPBSwH8RsouGCBcMBktLt1AymVo2TVsBVCY4b6TnZ/M= +github.com/bytedance/gopkg v0.1.3/go.mod h1:576VvJ+eJgyCzdjS+c4+77QF3p7ubbtiKARP3TxducM= +github.com/bytedance/sonic v1.15.0 h1:/PXeWFaR5ElNcVE84U0dOHjiMHQOwNIx3K4ymzh/uSE= +github.com/bytedance/sonic v1.15.0/go.mod h1:tFkWrPz0/CUCLEF4ri4UkHekCIcdnkqXw9VduqpJh0k= +github.com/bytedance/sonic/loader v0.5.0 h1:gXH3KVnatgY7loH5/TkeVyXPfESoqSBSBEiDd5VjlgE= +github.com/bytedance/sonic/loader v0.5.0/go.mod h1:AR4NYCk5DdzZizZ5djGqQ92eEhCCcdf5x77udYiSJRo= +github.com/cloudwego/base64x v0.1.6 h1:t11wG9AECkCDk5fMSoxmufanudBtJ+/HemLstXDLI2M= +github.com/cloudwego/base64x v0.1.6/go.mod h1:OFcloc187FXDaYHvrNIjxSe8ncn0OOM8gEHfghB2IPU= +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/gabriel-vasile/mimetype v1.4.12 h1:e9hWvmLYvtp846tLHam2o++qitpguFiYCKbn0w9jyqw= +github.com/gabriel-vasile/mimetype v1.4.12/go.mod h1:d+9Oxyo1wTzWdyVUPMmXFvp4F9tea18J8ufA774AB3s= +github.com/gin-contrib/sse v1.1.0 h1:n0w2GMuUpWDVp7qSpvze6fAu9iRxJY4Hmj6AmBOU05w= +github.com/gin-contrib/sse v1.1.0/go.mod h1:hxRZ5gVpWMT7Z0B0gSNYqqsSCNIJMjzvm6fqCz9vjwM= +github.com/gin-gonic/gin v1.12.0 h1:b3YAbrZtnf8N//yjKeU2+MQsh2mY5htkZidOM7O0wG8= +github.com/gin-gonic/gin v1.12.0/go.mod h1:VxccKfsSllpKshkBWgVgRniFFAzFb9csfngsqANjnLc= +github.com/go-playground/assert/v2 v2.2.0 h1:JvknZsQTYeFEAhQwI4qEt9cyV5ONwRHC+lYKSsYSR8s= +github.com/go-playground/assert/v2 v2.2.0/go.mod h1:VDjEfimB/XKnb+ZQfWdccd7VUvScMdVu0Titje2rxJ4= +github.com/go-playground/locales v0.14.1 h1:EWaQ/wswjilfKLTECiXz7Rh+3BjFhfDFKv/oXslEjJA= +github.com/go-playground/locales v0.14.1/go.mod h1:hxrqLVvrK65+Rwrd5Fc6F2O76J/NuW9t0sjnWqG1slY= +github.com/go-playground/universal-translator v0.18.1 h1:Bcnm0ZwsGyWbCzImXv+pAJnYK9S473LQFuzCbDbfSFY= +github.com/go-playground/universal-translator v0.18.1/go.mod h1:xekY+UJKNuX9WP91TpwSH2VMlDf28Uj24BCp08ZFTUY= +github.com/go-playground/validator/v10 v10.30.1 h1:f3zDSN/zOma+w6+1Wswgd9fLkdwy06ntQJp0BBvFG0w= +github.com/go-playground/validator/v10 v10.30.1/go.mod h1:oSuBIQzuJxL//3MelwSLD5hc2Tu889bF0Idm9Dg26cM= +github.com/goccy/go-json v0.10.5 h1:Fq85nIqj+gXn/S5ahsiTlK3TmC85qgirsdTP/+DeaC4= +github.com/goccy/go-json v0.10.5/go.mod h1:oq7eo15ShAhp70Anwd5lgX2pLfOS3QCiwU/PULtXL6M= +github.com/goccy/go-yaml v1.19.2 h1:PmFC1S6h8ljIz6gMRBopkjP1TVT7xuwrButHID66PoM= +github.com/goccy/go-yaml v1.19.2/go.mod h1:XBurs7gK8ATbW4ZPGKgcbrY1Br56PdM69F7LkFRi1kA= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= +github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= +github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM= +github.com/jackc/pgx/v5 v5.11.0 h1:IzBBtyK9AHqf98cctWFifYSci2hgQR/cd56wB4p+ogg= +github.com/jackc/pgx/v5 v5.11.0/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QIIM0Jn4= +github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo= +github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= +github.com/jinzhu/inflection v1.0.0 h1:K317FqzuhWc8YvSVlFMCCUb36O/S9MCKRDI7QkRKD/E= +github.com/jinzhu/inflection v1.0.0/go.mod h1:h+uFLlag+Qp1Va5pdKtLDYj+kHp5pxUVkryuEj+Srlc= +github.com/jinzhu/now v1.1.5 h1:/o9tlHleP7gOFmsnYNz3RGnqzefHA47wQpKrrdTIwXQ= +github.com/jinzhu/now v1.1.5/go.mod h1:d3SSVoowX0Lcu0IBviAWJpolVfI5UJVZZ7cO71lE/z8= +github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM= +github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo= +github.com/klauspost/cpuid/v2 v2.3.0 h1:S4CRMLnYUhGeDFDqkGriYKdfoFlDnMtqTiI/sFzhA9Y= +github.com/klauspost/cpuid/v2 v2.3.0/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= +github.com/leodido/go-urn v1.4.0 h1:WT9HwE9SGECu3lg4d/dIA+jxlljEa1/ffXKmRjqdmIQ= +github.com/leodido/go-urn v1.4.0/go.mod h1:bvxc+MVxLKB4z00jd1z+Dvzr47oO32F/QSNjSBOlFxI= +github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= +github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= +github.com/mattn/go-sqlite3 v1.14.22 h1:2gZY6PC6kBnID23Tichd1K+Z0oS6nE/XwU+Vz/5o4kU= +github.com/mattn/go-sqlite3 v1.14.22/go.mod h1:Uh1q+B4BYcTPb+yiD3kU8Ct7aC0hY9fxUwlHK0RXw+Y= +github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= +github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd h1:TRLaZ9cD/w8PVh93nsPXa1VrQ6jlwL5oN8l14QlcNfg= +github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= +github.com/modern-go/reflect2 v1.0.2 h1:xBagoLtFs94CBntxluKeaWgTMpvLxC4ur3nMaC9Gz0M= +github.com/modern-go/reflect2 v1.0.2/go.mod h1:yWuevngMOJpCy52FWWMvUC8ws7m/LJsjYzDa0/r8luk= +github.com/pelletier/go-toml/v2 v2.2.4 h1:mye9XuhQ6gvn5h28+VilKrrPoQVanw5PMw/TB0t5Ec4= +github.com/pelletier/go-toml/v2 v2.2.4/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8t9Stt+h56mCY= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/quic-go/qpack v0.6.0 h1:g7W+BMYynC1LbYLSqRt8PBg5Tgwxn214ZZR34VIOjz8= +github.com/quic-go/qpack v0.6.0/go.mod h1:lUpLKChi8njB4ty2bFLX2x4gzDqXwUpaO1DP9qMDZII= +github.com/quic-go/quic-go v0.59.0 h1:OLJkp1Mlm/aS7dpKgTc6cnpynnD2Xg7C1pwL6vy/SAw= +github.com/quic-go/quic-go v0.59.0/go.mod h1:upnsH4Ju1YkqpLXC305eW3yDZ4NfnNbmQRCMWS58IKU= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw= +github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpEOglKo= +github.com/stretchr/objx v0.5.2/go.mod h1:FRsXN1f5AsAjCGJKqEizvkpNtU+EGNCLh3NxZ/8L+MA= +github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= +github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU= +github.com/stretchr/testify v1.8.4/go.mod h1:sz/lmYIOXD/1dqDmKjjqLyZ2RngseejIcXlSw2iwfAo= +github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +github.com/twitchyliquid64/golang-asm v0.15.1 h1:SU5vSMR7hnwNxj24w34ZyCi/FmDZTkS4MhqMhdFk5YI= +github.com/twitchyliquid64/golang-asm v0.15.1/go.mod h1:a1lVb/DtPvCB8fslRZhAngC2+aY1QWCk3Cedj/Gdt08= +github.com/ugorji/go/codec v1.3.1 h1:waO7eEiFDwidsBN6agj1vJQ4AG7lh2yqXyOXqhgQuyY= +github.com/ugorji/go/codec v1.3.1/go.mod h1:pRBVtBSKl77K30Bv8R2P+cLSGaTtex6fsA2Wjqmfxj4= +go.mongodb.org/mongo-driver/v2 v2.5.0 h1:yXUhImUjjAInNcpTcAlPHiT7bIXhshCTL3jVBkF3xaE= +go.mongodb.org/mongo-driver/v2 v2.5.0/go.mod h1:yOI9kBsufol30iFsl1slpdq1I0eHPzybRWdyYUs8K/0= +go.uber.org/mock v0.6.0 h1:hyF9dfmbgIX5EfOdasqLsWD6xqpNZlXblLB/Dbnwv3Y= +go.uber.org/mock v0.6.0/go.mod h1:KiVJ4BqZJaMj4svdfmHM0AUx4NJYO8ZNpPnZn1Z+BBU= +golang.org/x/arch v0.22.0 h1:c/Zle32i5ttqRXjdLyyHZESLD/bB90DCU1g9l/0YBDI= +golang.org/x/arch v0.22.0/go.mod h1:dNHoOeKiyja7GTvF9NJS1l3Z2yntpQNzgrjh1cU103A= +golang.org/x/crypto v0.48.0 h1:/VRzVqiRSggnhY7gNRxPauEQ5Drw9haKdM0jqfcCFts= +golang.org/x/crypto v0.48.0/go.mod h1:r0kV5h3qnFPlQnBSrULhlsRfryS2pmewsg+XfMgkVos= +golang.org/x/net v0.51.0 h1:94R/GTO7mt3/4wIKpcR5gkGmRLOuE/2hNGeWq/GBIFo= +golang.org/x/net v0.51.0/go.mod h1:aamm+2QF5ogm02fjy5Bb7CQ0WMt1/WVM7FtyaTLlA9Y= +golang.org/x/sync v0.19.0 h1:vV+1eWNmZ5geRlYjzm2adRgW2/mcpevXNg50YZtPCE4= +golang.org/x/sync v0.19.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= +golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.41.0 h1:Ivj+2Cp/ylzLiEU89QhWblYnOE9zerudt9Ftecq2C6k= +golang.org/x/sys v0.41.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= +golang.org/x/text v0.34.0 h1:oL/Qq0Kdaqxa1KbNeMKwQq0reLCCaFtqu2eNuSeNHbk= +golang.org/x/text v0.34.0/go.mod h1:homfLqTYRFyVYemLBFl5GgL/DWEiH5wcsQ5gSh1yziA= +google.golang.org/protobuf v1.36.10 h1:AYd7cD/uASjIL6Q9LiTjz8JLcrh/88q5UObnmY3aOOE= +google.golang.org/protobuf v1.36.10/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gorm.io/driver/postgres v1.6.3 h1:bAn6O2pUa8LtpWEvL5NFU4+52Tfx8Ut7IVaIacCLcI0= +gorm.io/driver/postgres v1.6.3/go.mod h1:0c4fQA44XhOklXDkgtuKqysHCycTa5i9e3EIpDGCwXk= +gorm.io/driver/sqlite v1.6.0 h1:WHRRrIiulaPiPFmDcod6prc4l2VGVWHz80KspNsxSfQ= +gorm.io/driver/sqlite v1.6.0/go.mod h1:AO9V1qIQddBESngQUKWL9yoH93HIeA1X6V633rBwyT8= +gorm.io/gorm v1.31.2 h1:3o8FXNo9v9S858gil+3LlZA1LkCOzgb4g5BL64FgaCo= +gorm.io/gorm v1.31.2/go.mod h1:XyQVbO2k6YkOis7C2437jSit3SsDK72s7n7rsSHd+Gs= diff --git a/applications/task-dependencies/internal/board/cloud_test.go b/applications/task-dependencies/internal/board/cloud_test.go new file mode 100644 index 00000000..9a76414a --- /dev/null +++ b/applications/task-dependencies/internal/board/cloud_test.go @@ -0,0 +1,172 @@ +package board + +import ( + "context" + "crypto/x509" + "encoding/json" + "fmt" + "github.com/jackc/pgx/v5" + "gorm.io/gorm" + "net/http" + "os" + "strings" + "sync/atomic" + "testing" + "time" +) + +func cloud(t *testing.T) { + t.Helper() + if os.Getenv("TASKS_CLOUD_TESTS") != "1" { + t.Skip("requires disposable real Cloud fixture") + } +} +func TestCloudGORMTLS(t *testing.T) { + cloud(t) + ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second) + defer cancel() + config, err := pgx.ParseConfig("") + if err != nil { + t.Fatal(err) + } + db, pool, err := openConfig(ctx, config) + if err != nil { + t.Fatal("official CA positive failed") + } + defer pool.Close() + var ssl bool + if err = db.Raw("SELECT ssl FROM pg_stat_ssl WHERE pid=pg_backend_pid()").Scan(&ssl).Error; err != nil || !ssl { + t.Fatal("actual GORM connection has no observed TLS", err) + } + for _, control := range []string{"wrong CA", "wrong hostname"} { + changed := config.Copy() + changed.TLSConfig = config.TLSConfig.Clone() + if control == "wrong CA" { + changed.TLSConfig.RootCAs = x509.NewCertPool() + } else { + changed.TLSConfig.ServerName = "wrong-name.invalid" + } + _, failed, err := openConfig(ctx, changed) + if failed != nil { + failed.Close() + } + if err == nil { + t.Fatal(control, "unexpectedly connected") + } + message := err.Error() + if control == "wrong CA" && !strings.Contains(message, "unknown authority") { + t.Fatal("wrong CA did not fail specifically", err) + } + if control == "wrong hostname" && !strings.Contains(message, "not wrong-name.invalid") { + t.Fatal("wrong hostname did not fail specifically", err) + } + fmt.Println("Actual GORM/pgx connection rejected", control, "for certificate verification reason.") + } +} +func TestCloudCoherentGraphSnapshot(t *testing.T) { + cloud(t) + ctx, cancel := context.WithTimeout(context.Background(), 25*time.Second) + defer cancel() + db, pool, err := Open(ctx) + if err != nil { + t.Fatal(err) + } + defer pool.Close() + project := os.Getenv("PROBE_PROJECT_ID") + if _, err = UUID(project); err != nil { + t.Fatal("set PROBE_PROJECT_ID to dedicated owner-created fixture") + } + paused := make(chan struct{}) + release := make(chan struct{}) + var once atomic.Bool + var readerPID atomic.Int64 + err = db.Callback().Query().After("gorm:query").Register("acceptance:pause_project", func(tx *gorm.DB) { + if tx.Statement.Table == "projects" && once.CompareAndSwap(false, true) { + var pid int64 + if err := tx.Statement.ConnPool.QueryRowContext(ctx, "SELECT pg_backend_pid()").Scan(&pid); err != nil { + t.Error(err) + } + readerPID.Store(pid) + close(paused) + select { + case <-release: + case <-ctx.Done(): + } + } + }) + if err != nil { + t.Fatal(err) + } + type outcome struct { + snapshot Snapshot + err error + } + read := make(chan outcome, 1) + go func() { + snapshot, err := (Store{DB: db}).Snapshot(ctx, "10000000-0000-4000-8000-000000000001", project) + read <- outcome{snapshot, err} + }() + select { + case <-paused: + case <-ctx.Done(): + t.Fatal("project SELECT was not paused") + } + write := make(chan error, 1) + go func() { + request, err := http.NewRequestWithContext(ctx, "POST", "http://127.0.0.1:8090/api/projects/"+project+"/tasks", strings.NewReader(`{"title":"Coherent writer"}`)) + if err != nil { + write <- err + return + } + request.Header.Set("Authorization", "Bearer "+os.Getenv("TASKS_NORTH_TOKEN")) + request.Header.Set("Content-Type", "application/json") + response, err := http.DefaultClient.Do(request) + if err != nil { + write <- err + return + } + defer response.Body.Close() + if response.StatusCode != 201 { + write <- fmt.Errorf("writer status %d", response.StatusCode) + return + } + write <- nil + }() + deadline := time.Now().Add(8 * time.Second) + observed := false + for time.Now().Before(deadline) { + var count int + err = db.Raw("SELECT count(DISTINCT pid) FROM pg_locks WHERE NOT granted AND ? = ANY(pg_blocking_pids(pid))", readerPID.Load()).Scan(&count).Error + if err != nil { + t.Fatal(err) + } + if count > 0 { + observed = true + break + } + time.Sleep(50 * time.Millisecond) + } + close(release) + result := <-read + if !observed { + t.Fatal("actual HTTP writer was not observed waiting") + } + if result.err != nil { + t.Fatal(result.err) + } + if len(result.snapshot.Tasks) != 0 || result.snapshot.Project.Revision != 0 { + t.Fatal("reader mixed project/header with later task write") + } + if err = <-write; err != nil { + t.Fatal(err) + } + if err = db.Callback().Query().Remove("acceptance:pause_project"); err != nil { + t.Fatal(err) + } + after, err := (Store{DB: db}).Snapshot(ctx, "10000000-0000-4000-8000-000000000001", project) + if err != nil || after.Project.Revision != 1 || len(after.Tasks) != 1 { + t.Fatal("writer not visible in later complete snapshot", err) + } + encoded, _ := json.Marshal(after) + fmt.Println("Observed SHARE snapshot blocked independent HTTP writer; old header/graph agree, later header/graph agree:", string(encoded)) +} diff --git a/applications/task-dependencies/internal/board/database.go b/applications/task-dependencies/internal/board/database.go new file mode 100644 index 00000000..520801b0 --- /dev/null +++ b/applications/task-dependencies/internal/board/database.go @@ -0,0 +1,48 @@ +package board + +import ( + "context" + "database/sql" + "errors" + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/stdlib" + "gorm.io/driver/postgres" + "gorm.io/gorm" + "gorm.io/gorm/logger" + "os" + "time" +) + +func Open(ctx context.Context) (*gorm.DB, *sql.DB, error) { + if os.Getenv("PGSSLMODE") != "verify-full" || os.Getenv("PGSSLROOTCERT") == "" { + return nil, nil, errors.New("PGSSLMODE=verify-full and official PGSSLROOTCERT are required") + } + config, err := pgx.ParseConfig("") + if err != nil { + return nil, nil, err + } + config.ConnectTimeout = 5 * time.Second + config.RuntimeParams["application_name"] = "task-dependencies" + config.RuntimeParams["statement_timeout"] = "10000" + config.RuntimeParams["lock_timeout"] = "5000" + config.RuntimeParams["idle_in_transaction_session_timeout"] = "15000" + return openConfig(ctx, config) +} + +func openConfig(ctx context.Context, config *pgx.ConnConfig) (*gorm.DB, *sql.DB, error) { + pool := stdlib.OpenDB(*config) + pool.SetMaxOpenConns(4) + pool.SetMaxIdleConns(4) + pool.SetConnMaxIdleTime(30 * time.Second) + pool.SetConnMaxLifetime(15 * time.Minute) + db, err := gorm.Open(postgres.New(postgres.Config{Conn: pool}), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent), DisableAutomaticPing: true, SkipDefaultTransaction: true}) + if err != nil { + pool.Close() + return nil, nil, err + } + if err = pool.PingContext(ctx); err != nil { + pool.Close() + return nil, nil, err + } + return db, pool, nil +} diff --git a/applications/task-dependencies/internal/board/domain.go b/applications/task-dependencies/internal/board/domain.go new file mode 100644 index 00000000..61335301 --- /dev/null +++ b/applications/task-dependencies/internal/board/domain.go @@ -0,0 +1,122 @@ +package board + +import ( + "crypto/rand" + "encoding/hex" + "errors" + "strings" + "time" + "unicode" + "unicode/utf8" +) + +type Problem struct { + Status int + Code, Message string +} + +func (p *Problem) Error() string { return p.Message } +func conflict(code, message string) error { return &Problem{409, code, message} } + +var NotFound = &Problem{404, "NOT_FOUND", "Project or task not found."} +var BadInput = &Problem{422, "INVALID_INPUT", "Use a title of 1–100 characters without control or replacement characters."} + +func UUID(value string) (string, error) { + if len(value) != 36 || value[8] != '-' || value[13] != '-' || value[18] != '-' || value[23] != '-' { + return "", errors.New("invalid UUID") + } + raw := strings.ReplaceAll(value, "-", "") + if len(raw) != 32 { + return "", errors.New("invalid UUID separators") + } + decoded, err := hex.DecodeString(raw) + if err != nil || len(decoded) != 16 { + return "", errors.New("invalid UUID hex") + } + return strings.ToLower(value), nil +} +func NewUUID() (string, error) { + var bytes [16]byte + if _, err := rand.Read(bytes[:]); err != nil { + return "", err + } + bytes[6] = (bytes[6] & 15) | 64 + bytes[8] = (bytes[8] & 63) | 128 + s := hex.EncodeToString(bytes[:]) + return s[:8] + "-" + s[8:12] + "-" + s[12:16] + "-" + s[16:20] + "-" + s[20:], nil +} +func Title(value string) (string, error) { + if !utf8.ValidString(value) || len(value) > 400 { + return "", BadInput + } + for _, r := range value { + if unicode.IsControl(r) || unicode.Is(unicode.Cs, r) || r == utf8.RuneError { + return "", BadInput + } + } + value = strings.TrimSpace(value) + if n := utf8.RuneCountInString(value); n < 1 || n > 100 { + return "", BadInput + } + return value, nil +} + +type Project struct { + ID string `json:"id"` + AccountID string `json:"-"` + Name string `json:"name"` + Revision int64 `json:"revision"` + CreatedAt time.Time `json:"created_at"` +} + +func (Project) TableName() string { return "task_dependencies.projects" } + +type Task struct { + ID string `json:"id"` + ProjectID string `json:"project_id"` + Title string `json:"title"` + Done bool `json:"done"` + CreatedAt time.Time `json:"created_at"` + DoneAt *time.Time `json:"done_at"` +} + +func (Task) TableName() string { return "task_dependencies.tasks" } + +type Edge struct { + ProjectID string `json:"project_id"` + TaskID string `json:"task_id"` + PrerequisiteID string `json:"prerequisite_id"` + CreatedAt time.Time `json:"created_at"` +} + +func (Edge) TableName() string { return "task_dependencies.edges" } + +type Snapshot struct { + Project Project `json:"project"` + Tasks []Task `json:"tasks"` + Edges []Edge `json:"edges"` +} + +// Edges point from a task to its prerequisite. A new task -> prerequisite edge +// cycles exactly when the prerequisite already reaches the task. +func Reaches(edges []Edge, start, target string) bool { + next := make(map[string][]string) + for _, edge := range edges { + next[edge.TaskID] = append(next[edge.TaskID], edge.PrerequisiteID) + } + visited := make(map[string]bool) + pending := []string{start} + for len(pending) > 0 { + id := pending[len(pending)-1] + pending = pending[:len(pending)-1] + if id == target { + return true + } + if visited[id] { + continue + } + visited[id] = true + pending = append(pending, next[id]...) + } + return false +} diff --git a/applications/task-dependencies/internal/board/domain_test.go b/applications/task-dependencies/internal/board/domain_test.go new file mode 100644 index 00000000..bf93118c --- /dev/null +++ b/applications/task-dependencies/internal/board/domain_test.go @@ -0,0 +1,55 @@ +package board + +import ( + "encoding/json" + "strings" + "testing" +) + +func TestUUIDSeparators(t *testing.T) { + valid := "ABCDEF01-2345-6789-ABCD-EF0123456789" + actual, err := UUID(valid) + if err != nil || actual != strings.ToLower(valid) { + t.Fatal(actual, err) + } + // Required separators plus two additional hyphens still give length 36, + // but previously decoded only 15 bytes after removing every hyphen. + invalid := "--CDEF01-2345-6789-ABCD-EF0123456789" + if len(invalid) != 36 { + t.Fatal("fixture length") + } + if _, err = UUID(invalid); err == nil { + t.Fatal("extra separator accepted") + } +} +func TestTitleUnicodeBoundary(t *testing.T) { + for _, encoded := range []string{`"\ud800"`, `"bad\u007f"`, `"bad\u0085"`} { + var text string + if err := json.Unmarshal([]byte(encoded), &text); err != nil { + t.Fatal(err) + } + if _, err := Title(text); err == nil { + t.Fatal("unsafe title accepted") + } + } + var paired string + if err := json.Unmarshal([]byte(`"Plan \ud83d\ude80"`), &paired); err != nil { + t.Fatal(err) + } + if _, err := Title(paired); err != nil { + t.Fatal(err) + } + if title, err := Title(" Plan launch "); err != nil || title != "Plan launch" { + t.Fatal(title, err) + } +} +func TestReachabilityLongCycle(t *testing.T) { + edges := []Edge{{TaskID: "a", PrerequisiteID: "b"}, {TaskID: "b", PrerequisiteID: "c"}, {TaskID: "c", PrerequisiteID: "d"}} + if !Reaches(edges, "a", "d") || Reaches(edges, "d", "a") || Reaches(edges, "x", "a") { + t.Fatal("wrong path") + } + edges = append(edges, Edge{TaskID: "d", PrerequisiteID: "a"}) + if !Reaches(edges, "b", "a") { + t.Fatal("visited cycle traversal") + } +} diff --git a/applications/task-dependencies/internal/board/store.go b/applications/task-dependencies/internal/board/store.go new file mode 100644 index 00000000..7a5e22dd --- /dev/null +++ b/applications/task-dependencies/internal/board/store.go @@ -0,0 +1,220 @@ +package board + +import ( + "context" + "errors" + "gorm.io/gorm" + "gorm.io/gorm/clause" +) + +type Store struct{ DB *gorm.DB } + +func (s Store) Projects(ctx context.Context, account, cursor string) ([]Project, error) { + list := []Project{} + query := s.DB.WithContext(ctx).Where("account_id = ?", account) + if cursor != "" { + query = query.Where("id > ?", cursor) + } + err := query.Order("id ASC").Limit(20).Find(&list).Error + return list, err +} +func lockedProject(tx *gorm.DB, account, id, strength string) (Project, error) { + var project Project + result := tx.Clauses(clause.Locking{Strength: strength}).Where("account_id = ? AND id = ?", account, id).Take(&project) + if errors.Is(result.Error, gorm.ErrRecordNotFound) { + return project, NotFound + } + if result.Error != nil { + return project, result.Error + } + if result.RowsAffected != 1 { + return project, NotFound + } + return project, nil +} +func graph(tx *gorm.DB, project Project) (Snapshot, error) { + snapshot := Snapshot{Project: project, Tasks: []Task{}, Edges: []Edge{}} + if err := tx.Where("project_id = ?", project.ID).Order("id ASC").Limit(101).Find(&snapshot.Tasks).Error; err != nil { + return snapshot, err + } + if err := tx.Where("project_id = ?", project.ID).Order("task_id ASC, prerequisite_id ASC").Limit(301).Find(&snapshot.Edges).Error; err != nil { + return snapshot, err + } + if len(snapshot.Tasks) > 100 || len(snapshot.Edges) > 300 { + return snapshot, errors.New("stored graph exceeds bounds") + } + return snapshot, nil +} +func (s Store) Snapshot(ctx context.Context, account, id string) (Snapshot, error) { + var result Snapshot + err := s.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + project, err := lockedProject(tx, account, id, "SHARE") + if err != nil { + return err + } + result, err = graph(tx, project) + return err + }) + return result, err +} +func bump(tx *gorm.DB, project Project) error { + if project.Revision >= 1000000000 { + return conflict("REVISION_LIMIT", "This project has reached its revision limit.") + } + result := tx.Model(&Project{}).Where("id = ? AND revision = ?", project.ID, project.Revision).UpdateColumn("revision", gorm.Expr("revision + 1")) + if result.Error != nil { + return result.Error + } + if result.RowsAffected != 1 { + return errors.New("project revision update missing") + } + return nil +} +func (s Store) mutate(ctx context.Context, account, id string, operation func(*gorm.DB, Snapshot) error) error { + return s.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + project, err := lockedProject(tx, account, id, "UPDATE") + if err != nil { + return err + } + snapshot, err := graph(tx, project) + if err != nil { + return err + } + return operation(tx, snapshot) + }) +} +func (s Store) CreateTask(ctx context.Context, account, project, title string) (Task, error) { + task := Task{ProjectID: project, Title: title} + id, err := NewUUID() + if err != nil { + return task, err + } + task.ID = id + err = s.mutate(ctx, account, project, func(tx *gorm.DB, snapshot Snapshot) error { + if len(snapshot.Tasks) >= 100 { + return conflict("TASK_LIMIT", "This project already has 100 tasks.") + } + if err := bump(tx, snapshot.Project); err != nil { + return err + } + result := tx.Clauses(clause.Returning{}).Create(&task) + if result.Error != nil { + return result.Error + } + if result.RowsAffected != 1 { + return errors.New("task insert missing") + } + return nil + }) + return task, err +} +func tasksByID(snapshot Snapshot) map[string]Task { + tasks := make(map[string]Task) + for _, task := range snapshot.Tasks { + tasks[task.ID] = task + } + return tasks +} +func (s Store) AddEdge(ctx context.Context, account, project, task, prerequisite string) (Edge, error) { + edge := Edge{ProjectID: project, TaskID: task, PrerequisiteID: prerequisite} + err := s.mutate(ctx, account, project, func(tx *gorm.DB, snapshot Snapshot) error { + tasks := tasksByID(snapshot) + dependent, ok := tasks[task] + if !ok { + return NotFound + } + if _, ok := tasks[prerequisite]; !ok { + return NotFound + } + for _, existing := range snapshot.Edges { + if existing.TaskID == task && existing.PrerequisiteID == prerequisite { + edge = existing + return nil + } + } + if dependent.Done { + return conflict("TASK_DONE", "A completed task cannot gain a prerequisite.") + } + if task == prerequisite || Reaches(snapshot.Edges, prerequisite, task) { + return conflict("CYCLE", "This prerequisite would create a cycle.") + } + if len(snapshot.Edges) >= 300 { + return conflict("EDGE_LIMIT", "This project already has 300 prerequisites.") + } + if err := bump(tx, snapshot.Project); err != nil { + return err + } + result := tx.Clauses(clause.Returning{}).Create(&edge) + if result.Error != nil { + return result.Error + } + if result.RowsAffected != 1 { + return errors.New("edge insert missing") + } + return nil + }) + return edge, err +} +func (s Store) RemoveEdge(ctx context.Context, account, project, task, prerequisite string) error { + return s.mutate(ctx, account, project, func(tx *gorm.DB, snapshot Snapshot) error { + tasks := tasksByID(snapshot) + if _, ok := tasks[task]; !ok { + return NotFound + } + if _, ok := tasks[prerequisite]; !ok { + return NotFound + } + exists := false + for _, edge := range snapshot.Edges { + if edge.TaskID == task && edge.PrerequisiteID == prerequisite { + exists = true + break + } + } + if !exists { + return nil + } + if err := bump(tx, snapshot.Project); err != nil { + return err + } + result := tx.Where("project_id = ? AND task_id = ? AND prerequisite_id = ?", project, task, prerequisite).Delete(&Edge{}) + if result.Error != nil { + return result.Error + } + if result.RowsAffected != 1 { + return errors.New("edge removal missing") + } + return nil + }) +} +func (s Store) Complete(ctx context.Context, account, project, id string) (Task, error) { + var task Task + err := s.mutate(ctx, account, project, func(tx *gorm.DB, snapshot Snapshot) error { + tasks := tasksByID(snapshot) + var ok bool + task, ok = tasks[id] + if !ok { + return NotFound + } + if task.Done { + return nil + } + for _, edge := range snapshot.Edges { + if edge.TaskID == id && !tasks[edge.PrerequisiteID].Done { + return conflict("PREREQUISITES_OPEN", "Complete this task’s prerequisites first.") + } + } + if err := bump(tx, snapshot.Project); err != nil { + return err + } + result := tx.Model(&task).Clauses(clause.Returning{}).Where("project_id = ? AND id = ? AND done = false", project, id).Updates(map[string]interface{}{"done": true, "done_at": gorm.Expr("clock_timestamp()")}) + if result.Error != nil { + return result.Error + } + if result.RowsAffected != 1 { + return errors.New("task completion missing") + } + return nil + }) + return task, err +} diff --git a/applications/task-dependencies/internal/httpapi/api.go b/applications/task-dependencies/internal/httpapi/api.go new file mode 100644 index 00000000..24084a4a --- /dev/null +++ b/applications/task-dependencies/internal/httpapi/api.go @@ -0,0 +1,290 @@ +package httpapi + +import ( + "bytes" + "context" + "crypto/sha256" + "crypto/subtle" + "encoding/json" + "errors" + "github.com/ClickHouse/examples/applications/task-dependencies/internal/board" + "github.com/gin-gonic/gin" + "io" + "net/http" + "strings" + "time" +) + +const North = "10000000-0000-4000-8000-000000000001" +const South = "10000000-0000-4000-8000-000000000002" + +type Credential struct { + Digest [32]byte + Account string +} + +func Credentials(north, south string) ([]Credential, error) { + if len(north) < 32 || len(south) < 32 || len(north) > 256 || len(south) > 256 || north == south { + return nil, errors.New("two distinct bearer tokens of 32–256 bytes required") + } + return []Credential{{sha256.Sum256([]byte(north)), North}, {sha256.Sum256([]byte(south)), South}}, nil +} +func report(c *gin.Context, err error) { + var problem *board.Problem + if errors.As(err, &problem) { + c.JSON(problem.Status, gin.H{"error": problem.Code, "message": problem.Message}) + return + } + c.JSON(503, gin.H{"error": "DATABASE_UNAVAILABLE", "message": "The operation could not be confirmed. Read the current state before retrying a task creation."}) +} +func bad(c *gin.Context, message string) { + c.AbortWithStatusJSON(400, gin.H{"error": "BAD_REQUEST", "message": message}) +} +func id(c *gin.Context, key string) (string, bool) { + value, err := board.UUID(c.Param(key)) + if err != nil { + bad(c, "Use a valid UUID path.") + return "", false + } + return value, true +} +func decode(c *gin.Context, target any) bool { + if c.GetHeader("Content-Type") != "application/json" { + bad(c, "Use Content-Type: application/json.") + return false + } + c.Request.Body = http.MaxBytesReader(c.Writer, c.Request.Body, 8192) + body, err := io.ReadAll(c.Request.Body) + if err != nil { + bad(c, "JSON body must fit 8 KiB.") + return false + } + // Reject duplicate top-level properties before decoding the typed shape. + scan := json.NewDecoder(bytes.NewReader(body)) + token, err := scan.Token() + if err != nil || token != json.Delim('{') { + bad(c, "Use one JSON object.") + return false + } + seen := map[string]bool{} + for scan.More() { + key, err := scan.Token() + if err != nil { + bad(c, "Malformed JSON.") + return false + } + name, ok := key.(string) + if !ok || seen[name] { + bad(c, "Duplicate JSON property.") + return false + } + seen[name] = true + var raw json.RawMessage + if scan.Decode(&raw) != nil { + bad(c, "Malformed JSON.") + return false + } + } + if _, err = scan.Token(); err != nil { + bad(c, "Malformed JSON.") + return false + } + var extra any + if scan.Decode(&extra) != io.EOF { + bad(c, "Use exactly one JSON object.") + return false + } + decoder := json.NewDecoder(bytes.NewReader(body)) + decoder.DisallowUnknownFields() + if decoder.Decode(target) != nil { + bad(c, "Use the documented JSON fields and types.") + return false + } + return true +} +func Router(store board.Store, credentials []Credential) *gin.Engine { + gin.SetMode(gin.ReleaseMode) + r := gin.New() + r.SetTrustedProxies(nil) + r.Use(gin.CustomRecoveryWithWriter(io.Discard, func(c *gin.Context, _ any) { + c.AbortWithStatusJSON(503, gin.H{"error": "REQUEST_FAILED", "message": "The request could not be completed."}) + })) + r.Use(func(c *gin.Context) { + c.Header("Cache-Control", "no-store") + c.Header("X-Content-Type-Options", "nosniff") + ctx, cancel := context.WithTimeout(c.Request.Context(), 10*time.Second) + defer cancel() + c.Request = c.Request.WithContext(ctx) + c.Next() + }) + r.GET("/healthz", func(c *gin.Context) { + pool, err := store.DB.DB() + if err != nil || pool.PingContext(c.Request.Context()) != nil { + c.Status(503) + return + } + c.JSON(200, gin.H{"status": "ready"}) + }) + api := r.Group("/api", func(c *gin.Context) { + header := c.GetHeader("Authorization") + if !strings.HasPrefix(header, "Bearer ") || len(header) > 263 { + c.AbortWithStatusJSON(401, gin.H{"error": "UNAUTHORIZED", "message": "Use an account bearer token."}) + return + } + digest := sha256.Sum256([]byte(strings.TrimPrefix(header, "Bearer "))) + account := "" + for _, credential := range credentials { + if subtle.ConstantTimeCompare(digest[:], credential.Digest[:]) == 1 { + account = credential.Account + } + } + if account == "" { + c.AbortWithStatusJSON(401, gin.H{"error": "UNAUTHORIZED", "message": "Use an account bearer token."}) + return + } + // Bearer headers are not ambient browser cookies. Reject foreign unsafe + // browser origins as an additional boundary; no cross-origin API is enabled. + if c.Request.Method != "GET" && c.Request.Method != "HEAD" { + origin := c.GetHeader("Origin") + if origin != "" && origin != "http://127.0.0.1:8090" { + c.AbortWithStatusJSON(403, gin.H{"error": "ORIGIN_REJECTED", "message": "Use the local API origin."}) + return + } + } + c.Set("account", account) + c.Next() + }) + api.GET("/projects", func(c *gin.Context) { + cursor := c.Query("after") + for key, values := range c.Request.URL.Query() { + if key != "after" || len(values) != 1 { + bad(c, "Only one after cursor is supported.") + return + } + } + if cursor != "" { + var err error + cursor, err = board.UUID(cursor) + if err != nil { + bad(c, "Use a valid UUID cursor.") + return + } + } + result, err := store.Projects(c.Request.Context(), c.GetString("account"), cursor) + if err != nil { + report(c, err) + return + } + c.JSON(200, result) + }) + api.GET("/projects/:project", func(c *gin.Context) { + project, ok := id(c, "project") + if !ok { + return + } + result, err := store.Snapshot(c.Request.Context(), c.GetString("account"), project) + if err != nil { + report(c, err) + return + } + c.JSON(200, result) + }) + api.POST("/projects/:project/tasks", func(c *gin.Context) { + project, ok := id(c, "project") + if !ok { + return + } + var input struct { + Title *string `json:"title"` + } + if !decode(c, &input) { + return + } + if input.Title == nil { + report(c, board.BadInput) + return + } + title, err := board.Title(*input.Title) + if err != nil { + report(c, err) + return + } + task, err := store.CreateTask(c.Request.Context(), c.GetString("account"), project, title) + if err != nil { + report(c, err) + return + } + c.JSON(201, task) + }) + api.POST("/projects/:project/tasks/:task/prerequisites", func(c *gin.Context) { + project, ok := id(c, "project") + if !ok { + return + } + task, ok := id(c, "task") + if !ok { + return + } + var input struct { + PrerequisiteID *string `json:"prerequisite_id"` + } + if !decode(c, &input) { + return + } + if input.PrerequisiteID == nil { + bad(c, "Use a prerequisite UUID.") + return + } + prerequisite, err := board.UUID(*input.PrerequisiteID) + if err != nil { + bad(c, "Use a prerequisite UUID.") + return + } + edge, err := store.AddEdge(c.Request.Context(), c.GetString("account"), project, task, prerequisite) + if err != nil { + report(c, err) + return + } + c.JSON(200, edge) + }) + api.DELETE("/projects/:project/tasks/:task/prerequisites/:prerequisite", func(c *gin.Context) { + project, ok := id(c, "project") + if !ok { + return + } + task, ok := id(c, "task") + if !ok { + return + } + prerequisite, ok := id(c, "prerequisite") + if !ok { + return + } + if err := store.RemoveEdge(c.Request.Context(), c.GetString("account"), project, task, prerequisite); err != nil { + report(c, err) + return + } + c.Status(204) + }) + api.POST("/projects/:project/tasks/:task/complete", func(c *gin.Context) { + project, ok := id(c, "project") + if !ok { + return + } + task, ok := id(c, "task") + if !ok { + return + } + var input struct{} + if !decode(c, &input) { + return + } + result, err := store.Complete(c.Request.Context(), c.GetString("account"), project, task) + if err != nil { + report(c, err) + return + } + c.JSON(200, result) + }) + return r +} diff --git a/applications/task-dependencies/internal/httpapi/api_test.go b/applications/task-dependencies/internal/httpapi/api_test.go new file mode 100644 index 00000000..9f033814 --- /dev/null +++ b/applications/task-dependencies/internal/httpapi/api_test.go @@ -0,0 +1,53 @@ +package httpapi + +import ( + "bytes" + "github.com/ClickHouse/examples/applications/task-dependencies/internal/board" + "net/http/httptest" + "strings" + "testing" +) + +func request(t *testing.T, path, body, origin, token string) int { + t.Helper() + credentials, err := Credentials(strings.Repeat("n", 32), strings.Repeat("s", 32)) + if err != nil { + t.Fatal(err) + } + r := Router(board.Store{}, credentials) + req := httptest.NewRequest("POST", path, bytes.NewBufferString(body)) + req.Header.Set("Content-Type", "application/json") + if token != "" { + req.Header.Set("Authorization", "Bearer "+token) + } + if origin != "" { + req.Header.Set("Origin", origin) + } + out := httptest.NewRecorder() + r.ServeHTTP(out, req) + return out.Code +} +func TestCredentialMapping(t *testing.T) { + if _, err := Credentials("short", "short"); err == nil { + t.Fatal("weak tokens accepted") + } + if status := request(t, "/api/projects/invalid/tasks", `{}`, "", ""); status != 401 { + t.Fatal(status) + } +} +func TestForeignOrigin(t *testing.T) { + if status := request(t, "/api/projects/20000000-0000-4000-8000-000000000001/tasks", `{"title":"Plan"}`, "https://foreign.invalid", strings.Repeat("n", 32)); status != 403 { + t.Fatal(status) + } +} +func TestStrictJSONAndUUID(t *testing.T) { + base := "/api/projects/20000000-0000-4000-8000-000000000001/tasks" + for _, body := range []string{`{"title":"first","title":"second"}`, `{"unknown":"x"}`, `{"title":42}`, `{"title":"x"} {}`, `[`} { + if status := request(t, base, body, "", strings.Repeat("n", 32)); status != 400 { + t.Fatal(status, body) + } + } + if status := request(t, "/api/projects/--000000-0000-4000-8000-000000000001/tasks", `{"title":"x"}`, "", strings.Repeat("n", 32)); status != 400 { + t.Fatal(status) + } +} diff --git a/applications/task-dependencies/scripts/migrate.sh b/applications/task-dependencies/scripts/migrate.sh new file mode 100644 index 00000000..f386cdce --- /dev/null +++ b/applications/task-dependencies/scripts/migrate.sh @@ -0,0 +1,48 @@ +#!/usr/bin/env bash +set -euo pipefail +cd "$(dirname "$0")/.." +if [[ "${PGUSER:-}" != tasks_owner ]]; then + printf '%s\n' 'Migrations must use tasks_owner.' >&2 + exit 1 +fi +if [[ "${PGSSLMODE:-}" != verify-full || -z "${PGSSLROOTCERT:-}" ]]; then + printf '%s\n' 'Set PGSSLMODE=verify-full and the official PGSSLROOTCERT.' >&2 + exit 1 +fi +direction="${1:-up}" +case "$direction" in + up) + checksum=$(sha256sum sql/migrations/001-up.sql | cut -d ' ' -f 1) + psql -X -v ON_ERROR_STOP=1 -v checksum="$checksum" <<'SQL' +BEGIN; +SELECT pg_advisory_xact_lock(5870041002); +CREATE TABLE IF NOT EXISTS task_dependencies.schema_migrations ( + version integer PRIMARY KEY, checksum text NOT NULL +); +SELECT NOT EXISTS (SELECT FROM task_dependencies.schema_migrations WHERE version = 1) AS needed \gset +\if :needed +\i sql/migrations/001-up.sql +INSERT INTO task_dependencies.schema_migrations VALUES (1, :'checksum'); +\else +SELECT checksum = :'checksum' AS agrees FROM task_dependencies.schema_migrations WHERE version = 1 \gset +\if :agrees +\echo Migration 1 already applied with matching checksum. +\else +\echo Migration checksum changed; refusing to continue. +\quit 1 +\endif +\endif +COMMIT; +SQL + ;; + down) + psql -X -v ON_ERROR_STOP=1 <<'SQL' +BEGIN; +SELECT pg_advisory_xact_lock(5870041002); +\i sql/migrations/001-down.sql +DELETE FROM task_dependencies.schema_migrations WHERE version = 1; +COMMIT; +SQL + ;; + *) printf '%s\n' 'Usage: scripts/migrate.sh up|down' >&2; exit 1 ;; +esac diff --git a/applications/task-dependencies/sql/bootstrap.sql b/applications/task-dependencies/sql/bootstrap.sql new file mode 100644 index 00000000..921e866b --- /dev/null +++ b/applications/task-dependencies/sql/bootstrap.sql @@ -0,0 +1,14 @@ +\set ON_ERROR_STOP on +\getenv owner_password TASKS_OWNER_PASSWORD +\getenv runtime_password TASKS_RUNTIME_PASSWORD +SELECT 'CREATE ROLE tasks_owner LOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE' WHERE NOT EXISTS(SELECT 1 FROM pg_roles WHERE rolname='tasks_owner') \gexec +SELECT 'CREATE ROLE tasks_runtime LOGIN NOINHERIT NOSUPERUSER NOCREATEDB NOCREATEROLE' WHERE NOT EXISTS(SELECT 1 FROM pg_roles WHERE rolname='tasks_runtime') \gexec +ALTER ROLE tasks_owner PASSWORD :'owner_password'; +ALTER ROLE tasks_runtime PASSWORD :'runtime_password'; +CREATE SCHEMA IF NOT EXISTS task_dependencies AUTHORIZATION tasks_owner; +SELECT format('REVOKE CREATE, TEMP ON DATABASE %I FROM PUBLIC',current_database()) \gexec +SELECT format('GRANT CONNECT ON DATABASE %I TO tasks_owner, tasks_runtime',current_database()) \gexec +REVOKE CREATE ON SCHEMA public FROM PUBLIC; +ALTER ROLE tasks_runtime SET statement_timeout='10s'; +ALTER ROLE tasks_runtime SET lock_timeout='5s'; +ALTER ROLE tasks_runtime SET idle_in_transaction_session_timeout='15s'; diff --git a/applications/task-dependencies/sql/grants.sql b/applications/task-dependencies/sql/grants.sql new file mode 100644 index 00000000..a0e564a8 --- /dev/null +++ b/applications/task-dependencies/sql/grants.sql @@ -0,0 +1,11 @@ +\set ON_ERROR_STOP on +REVOKE ALL ON SCHEMA task_dependencies FROM PUBLIC; +REVOKE ALL ON ALL TABLES IN SCHEMA task_dependencies FROM PUBLIC; +REVOKE ALL ON ALL FUNCTIONS IN SCHEMA task_dependencies FROM PUBLIC; +GRANT USAGE ON SCHEMA task_dependencies TO tasks_runtime; +GRANT SELECT ON task_dependencies.projects,task_dependencies.tasks,task_dependencies.edges,task_dependencies.schema_migrations TO tasks_runtime; +GRANT UPDATE(revision) ON task_dependencies.projects TO tasks_runtime; +GRANT INSERT ON task_dependencies.tasks,task_dependencies.edges TO tasks_runtime; +GRANT UPDATE(done,done_at) ON task_dependencies.tasks TO tasks_runtime; +GRANT DELETE ON task_dependencies.edges TO tasks_runtime; +GRANT EXECUTE ON FUNCTION task_dependencies.enforce_bounds() TO tasks_runtime; diff --git a/applications/task-dependencies/sql/migrations/001-down.sql b/applications/task-dependencies/sql/migrations/001-down.sql new file mode 100644 index 00000000..37baaf69 --- /dev/null +++ b/applications/task-dependencies/sql/migrations/001-down.sql @@ -0,0 +1,4 @@ +DROP TABLE task_dependencies.edges; +DROP TABLE task_dependencies.tasks; +DROP TABLE task_dependencies.projects; +DROP FUNCTION task_dependencies.enforce_bounds(); diff --git a/applications/task-dependencies/sql/migrations/001-up.sql b/applications/task-dependencies/sql/migrations/001-up.sql new file mode 100644 index 00000000..eb1e304c --- /dev/null +++ b/applications/task-dependencies/sql/migrations/001-up.sql @@ -0,0 +1,44 @@ +CREATE TABLE task_dependencies.projects ( + id uuid PRIMARY KEY, + account_id uuid NOT NULL, + name text NOT NULL CHECK (char_length(name) BETWEEN 1 AND 100), + revision bigint NOT NULL DEFAULT 0 CHECK (revision BETWEEN 0 AND 1000000000), + created_at timestamptz NOT NULL DEFAULT clock_timestamp() +); +CREATE INDEX projects_account_cursor ON task_dependencies.projects(account_id,id); +CREATE TABLE task_dependencies.tasks ( + id uuid PRIMARY KEY, + project_id uuid NOT NULL REFERENCES task_dependencies.projects(id), + title text NOT NULL CHECK (char_length(title) BETWEEN 1 AND 100), + done boolean NOT NULL DEFAULT false, + done_at timestamptz, + created_at timestamptz NOT NULL DEFAULT clock_timestamp(), + UNIQUE(project_id,id), + CHECK ((done AND done_at IS NOT NULL) OR (NOT done AND done_at IS NULL)) +); +CREATE TABLE task_dependencies.edges ( + project_id uuid NOT NULL REFERENCES task_dependencies.projects(id), + task_id uuid NOT NULL, + prerequisite_id uuid NOT NULL, + created_at timestamptz NOT NULL DEFAULT clock_timestamp(), + PRIMARY KEY(project_id,task_id,prerequisite_id), + FOREIGN KEY(project_id,task_id) REFERENCES task_dependencies.tasks(project_id,id), + FOREIGN KEY(project_id,prerequisite_id) REFERENCES task_dependencies.tasks(project_id,id), + CHECK (task_id <> prerequisite_id) +); +CREATE FUNCTION task_dependencies.enforce_bounds() RETURNS trigger +LANGUAGE plpgsql SECURITY INVOKER SET search_path=pg_catalog AS $$ +DECLARE count_now integer; +BEGIN + PERFORM 1 FROM task_dependencies.projects WHERE id=NEW.project_id FOR UPDATE; + IF TG_TABLE_NAME='tasks' THEN + SELECT count(*) INTO count_now FROM task_dependencies.tasks WHERE project_id=NEW.project_id; + IF count_now>=100 THEN RAISE EXCEPTION 'task limit' USING ERRCODE='23514'; END IF; + ELSE + SELECT count(*) INTO count_now FROM task_dependencies.edges WHERE project_id=NEW.project_id; + IF count_now>=300 THEN RAISE EXCEPTION 'edge limit' USING ERRCODE='23514'; END IF; + END IF; + RETURN NEW; +END $$; +CREATE TRIGGER task_bound BEFORE INSERT ON task_dependencies.tasks FOR EACH ROW EXECUTE FUNCTION task_dependencies.enforce_bounds(); +CREATE TRIGGER edge_bound BEFORE INSERT ON task_dependencies.edges FOR EACH ROW EXECUTE FUNCTION task_dependencies.enforce_bounds(); diff --git a/applications/task-dependencies/sql/seed.sql b/applications/task-dependencies/sql/seed.sql new file mode 100644 index 00000000..49d0e5b9 --- /dev/null +++ b/applications/task-dependencies/sql/seed.sql @@ -0,0 +1,11 @@ +\set ON_ERROR_STOP on +INSERT INTO task_dependencies.projects(id,account_id,name) VALUES + ('20000000-0000-4000-8000-000000000001','10000000-0000-4000-8000-000000000001','Release preparation'), + ('20000000-0000-4000-8000-000000000002','10000000-0000-4000-8000-000000000001','Office move'), + ('20000000-0000-4000-8000-000000000003','10000000-0000-4000-8000-000000000002','Workshop opening') ON CONFLICT(id) DO NOTHING; +INSERT INTO task_dependencies.tasks(id,project_id,title) VALUES + ('30000000-0000-4000-8000-000000000001','20000000-0000-4000-8000-000000000001','Review release notes'), + ('30000000-0000-4000-8000-000000000002','20000000-0000-4000-8000-000000000001','Publish release'), + ('30000000-0000-4000-8000-000000000003','20000000-0000-4000-8000-000000000003','Check room') ON CONFLICT(id) DO NOTHING; +INSERT INTO task_dependencies.edges(project_id,task_id,prerequisite_id) VALUES + ('20000000-0000-4000-8000-000000000001','30000000-0000-4000-8000-000000000002','30000000-0000-4000-8000-000000000001') ON CONFLICT DO NOTHING; diff --git a/applications/task-dependencies/tests/acceptance.py b/applications/task-dependencies/tests/acceptance.py new file mode 100644 index 00000000..43064301 --- /dev/null +++ b/applications/task-dependencies/tests/acceptance.py @@ -0,0 +1,377 @@ +"""Actual disposable Cloud/HTTP acceptance. No database substitutes.""" + +import concurrent.futures +import contextlib +import json +import os +from pathlib import Path +import threading +import time +import unittest +import urllib.error +import urllib.request +import uuid +import psycopg + +BASE = os.environ.get("API_BASE_URL", "http://127.0.0.1:8090") +NORTH = "10000000-0000-4000-8000-000000000001" +SOUTH = "10000000-0000-4000-8000-000000000002" +TOKEN = os.environ["TASKS_NORTH_TOKEN"] +SOUTH_TOKEN = os.environ["TASKS_SOUTH_TOKEN"] + + +def owner(): + return psycopg.connect( + user=os.environ["TEST_OWNER_USER"], + password=os.environ["TEST_OWNER_PASSWORD"], + autocommit=True, + ) + + +def call(path, token=TOKEN, method="GET", payload=None, raw=None, origin=None): + body = ( + raw + if raw is not None + else (json.dumps(payload).encode() if payload is not None else None) + ) + headers = {"Authorization": "Bearer " + token, "Content-Type": "application/json"} + if origin: + headers["Origin"] = origin + request = urllib.request.Request( + BASE + "/api" + path, data=body, headers=headers, method=method + ) + try: + response = urllib.request.urlopen(request, timeout=20) + except urllib.error.HTTPError as error: + response = error + with response: + data = response.read() + return response.status, json.loads(data) if data else None + + +def observe_waiters(connection, blocker, expected): + deadline = time.monotonic() + 8 + while time.monotonic() < deadline: + count = connection.execute( + """WITH RECURSIVE waiting(pid) AS ( + SELECT pid FROM pg_locks WHERE NOT granted AND %s = ANY(pg_blocking_pids(pid)) + UNION SELECT l.pid FROM pg_locks l JOIN waiting w ON w.pid=ANY(pg_blocking_pids(l.pid)) WHERE NOT l.granted + ) SELECT count(DISTINCT pid) FROM waiting""", + (blocker,), + ).fetchone()[0] + if count >= expected: + print( + f"Observed {count} independent HTTP transactions in the real project lock queue.", + flush=True, + ) + return + time.sleep(0.05) + raise AssertionError( + "Expected HTTP transactions were not observed waiting on project lock." + ) + + +class CloudHTTP(unittest.TestCase): + def setUp(self): + self.project = str(uuid.uuid4()) + self.other = str(uuid.uuid4()) + with owner() as connection: + connection.execute( + "INSERT INTO task_dependencies.projects(id,account_id,name) VALUES (%s,%s,'Acceptance north'),(%s,%s,'Acceptance south')", + (self.project, NORTH, self.other, SOUTH), + ) + + def tearDown(self): + with owner() as connection: + connection.execute( + "DELETE FROM task_dependencies.edges WHERE project_id=ANY(%s::uuid[])", + ([self.project, self.other],), + ) + connection.execute( + "DELETE FROM task_dependencies.tasks WHERE project_id=ANY(%s::uuid[])", + ([self.project, self.other],), + ) + connection.execute( + "DELETE FROM task_dependencies.projects WHERE id=ANY(%s::uuid[])", + ([self.project, self.other],), + ) + + def task(self, title="Acceptance task", project=None, token=TOKEN): + status, result = call( + f"/projects/{project or self.project}/tasks", + token, + "POST", + {"title": title}, + ) + self.assertEqual(status, 201, result) + return result + + def edge(self, task, prerequisite): + return call( + f"/projects/{self.project}/tasks/{task}/prerequisites", + method="POST", + payload={"prerequisite_id": prerequisite}, + ) + + def complete(self, task): + return call( + f"/projects/{self.project}/tasks/{task}/complete", method="POST", payload={} + ) + + def snapshot(self): + status, result = call(f"/projects/{self.project}") + self.assertEqual(status, 200, result) + return result + + def test_01_workflow_duplicate_and_repeat_safe(self): + a = self.task("Review")["id"] + b = self.task("Publish")["id"] + first = self.edge(b, a) + self.assertEqual(first[0], 200) + before = self.snapshot()["project"]["revision"] + self.assertEqual(self.edge(b, a), first) + self.assertEqual(self.snapshot()["project"]["revision"], before) + self.assertEqual(self.complete(b)[0], 409) + self.assertTrue(self.complete(a)[1]["done"]) + result = self.complete(b) + self.assertEqual(result[0], 200) + self.assertTrue(result[1]["done"]) + self.assertIsNotNone(result[1]["done_at"]) + self.assertEqual(self.complete(b), result) + self.assertEqual(self.edge(b, a), first) + extra = self.task()["id"] + self.assertEqual(self.edge(b, extra)[0], 409) + path = f"/projects/{self.project}/tasks/{b}/prerequisites/{a}" + self.assertEqual(call(path, method="DELETE")[0], 204) + revision = self.snapshot()["project"]["revision"] + self.assertEqual(call(path, method="DELETE")[0], 204) + self.assertEqual(self.snapshot()["project"]["revision"], revision) + + def test_02_direct_and_long_cycles(self): + tasks = [self.task(str(i))["id"] for i in range(6)] + for a, b in zip(tasks, tasks[1:]): + self.assertEqual(self.edge(a, b)[0], 200) + self.assertEqual(self.edge(tasks[-1], tasks[0])[1]["error"], "CYCLE") + self.assertEqual(self.edge(tasks[0], tasks[0])[1]["error"], "CYCLE") + self.assertEqual(len(self.snapshot()["edges"]), 5) + + def test_03_opposite_edge_contention(self): + a = self.task()["id"] + b = self.task()["id"] + barrier = threading.Barrier(3) + with owner() as held, owner() as observer, concurrent.futures.ThreadPoolExecutor( + 2 + ) as pool: + with held.transaction(): + blocker = held.execute("SELECT pg_backend_pid()").fetchone()[0] + held.execute( + "SELECT id FROM task_dependencies.projects WHERE id=%s FOR UPDATE", + (self.project,), + ) + + def add(left, right): + barrier.wait() + return self.edge(left, right) + + pending = [pool.submit(add, a, b), pool.submit(add, b, a)] + barrier.wait() + observe_waiters(observer, blocker, 2) + results = [future.result() for future in pending] + self.assertEqual(sorted(result[0] for result in results), [200, 409]) + self.assertEqual(len(self.snapshot()["edges"]), 1) + + def test_04_completion_against_new_prerequisite(self): + a = self.task()["id"] + b = self.task()["id"] + barrier = threading.Barrier(3) + with owner() as held, owner() as observer, concurrent.futures.ThreadPoolExecutor( + 2 + ) as pool: + with held.transaction(): + blocker = held.execute("SELECT pg_backend_pid()").fetchone()[0] + held.execute( + "SELECT id FROM task_dependencies.projects WHERE id=%s FOR UPDATE", + (self.project,), + ) + + def add(): + barrier.wait() + return self.edge(a, b) + + def complete(): + barrier.wait() + return self.complete(a) + + pending = [pool.submit(add), pool.submit(complete)] + barrier.wait() + observe_waiters(observer, blocker, 2) + added, completed = [future.result() for future in pending] + snapshot = self.snapshot() + task = next(task for task in snapshot["tasks"] if task["id"] == a) + if added[0] == 200: + self.assertEqual(completed[0], 409) + self.assertFalse(task["done"]) + self.assertEqual(len(snapshot["edges"]), 1) + else: + self.assertEqual(added[0], 409) + self.assertEqual(completed[0], 200) + self.assertTrue(task["done"]) + self.assertEqual(snapshot["edges"], []) + print( + "Actual coordinated edge/complete outcome:", + added[0], + completed[0], + flush=True, + ) + + def test_05_multiwrite_rollback(self): + revision = self.snapshot()["project"]["revision"] + with owner() as connection: + connection.execute( + """CREATE FUNCTION task_dependencies.acceptance_fail() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN IF NEW.title='fail-after-revision' THEN RAISE EXCEPTION 'private owner fixture' USING ERRCODE='23514'; END IF; RETURN NEW; END $$""" + ) + connection.execute( + "CREATE TRIGGER acceptance_fail BEFORE INSERT ON task_dependencies.tasks FOR EACH ROW EXECUTE FUNCTION task_dependencies.acceptance_fail()" + ) + try: + status, result = call( + f"/projects/{self.project}/tasks", + method="POST", + payload={"title": "fail-after-revision"}, + ) + self.assertEqual(status, 503) + self.assertNotIn("private", json.dumps(result)) + after = self.snapshot() + self.assertEqual(after["project"]["revision"], revision) + self.assertEqual(after["tasks"], []) + finally: + with owner() as connection: + connection.execute( + "DROP TRIGGER acceptance_fail ON task_dependencies.tasks" + ) + connection.execute("DROP FUNCTION task_dependencies.acceptance_fail()") + + def test_06_account_scope_and_composite_keys(self): + a = self.task()["id"] + b = self.task(project=self.other, token=SOUTH_TOKEN)["id"] + self.assertEqual(call(f"/projects/{self.other}")[0], 404) + self.assertEqual( + call( + f"/projects/{self.other}/tasks", + method="POST", + payload={"title": "forged"}, + )[0], + 404, + ) + self.assertEqual(self.edge(a, b)[0], 404) + with owner() as connection: + for task, prerequisite, code in [(a, b, "23503"), (a, a, "23514")]: + with self.assertRaises(psycopg.Error) as error: + connection.execute( + "INSERT INTO task_dependencies.edges(project_id,task_id,prerequisite_id) VALUES (%s,%s,%s)", + (self.project, task, prerequisite), + ) + self.assertEqual(error.exception.sqlstate, code) + status, projects = call("/projects") + self.assertEqual(status, 200) + self.assertLessEqual(len(projects), 20) + self.assertEqual([p["id"] for p in projects], sorted(p["id"] for p in projects)) + self.assertNotIn(self.other, [p["id"] for p in projects]) + + def test_07_graph_bounds(self): + ids = [str(uuid.uuid4()) for _ in range(100)] + with owner() as connection: + with connection.transaction(): + for task in ids: + connection.execute( + "INSERT INTO task_dependencies.tasks(id,project_id,title) VALUES (%s,%s,%s)", + (task, self.project, "Bound fixture"), + ) + pairs = [(a, b) for i, a in enumerate(ids) for b in ids[i + 1 :]][:300] + for a, b in pairs: + connection.execute( + "INSERT INTO task_dependencies.edges(project_id,task_id,prerequisite_id) VALUES (%s,%s,%s)", + (self.project, a, b), + ) + with self.assertRaises(psycopg.Error) as failure: + connection.execute( + "INSERT INTO task_dependencies.tasks(id,project_id,title) VALUES (%s,%s,%s)", + (str(uuid.uuid4()), self.project, "One too many"), + ) + self.assertEqual(failure.exception.sqlstate, "23514") + with self.assertRaises(psycopg.Error) as failure: + connection.execute( + "INSERT INTO task_dependencies.edges(project_id,task_id,prerequisite_id) VALUES (%s,%s,%s)", + (self.project, ids[10], ids[99]), + ) + self.assertEqual(failure.exception.sqlstate, "23514") + self.assertEqual( + call( + f"/projects/{self.project}/tasks", + method="POST", + payload={"title": "Over limit"}, + )[1]["error"], + "TASK_LIMIT", + ) + self.assertEqual(self.edge(ids[10], ids[99])[1]["error"], "EDGE_LIMIT") + snapshot = self.snapshot() + self.assertEqual(len(snapshot["tasks"]), 100) + self.assertEqual(len(snapshot["edges"]), 300) + + def test_08_http_boundaries(self): + path = f"/projects/{self.project}/tasks" + for raw, status in [ + (b'{"title":"x","title":"y"}', 400), + (b'{"title":null}', 422), + (b'{"title":"\\ud800"}', 422), + (b'{"title":"bad\\u0085"}', 422), + (b'{"title":42}', 400), + (b'{"title":"' + b"x" * 9000 + b'"}', 400), + ]: + self.assertEqual(call(path, method="POST", raw=raw)[0], status) + self.assertEqual( + call(path, method="POST", payload={"title": "Plan 🚀"})[0], 201 + ) + self.assertEqual(call("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/projects/--000000-0000-4000-8000-000000000001")[0], 400) + self.assertEqual( + call( + path, + method="POST", + payload={"title": "x"}, + origin="https://foreign.invalid", + )[0], + 403, + ) + self.assertEqual(call("/projects", token="forged")[0], 401) + self.assertEqual(call("/projects?after=bad")[0], 400) + with owner() as connection: + connection.execute( + "UPDATE task_dependencies.projects SET revision=1000000000 WHERE id=%s", + (self.project,), + ) + self.assertEqual( + call(path, method="POST", payload={"title": "At revision cap"})[1]["error"], + "REVISION_LIMIT", + ) + + def test_09_runtime_roles(self): + with psycopg.connect() as connection: + flags = connection.execute( + "SELECT rolsuper,rolcreatedb,rolcreaterole FROM pg_roles WHERE rolname=current_user" + ).fetchone() + self.assertEqual(flags, (False, False, False)) + for sql in [ + "CREATE TABLE task_dependencies.forbidden(id int)", + "CREATE TEMP TABLE forbidden(id int)", + "SET ROLE tasks_owner", + "UPDATE task_dependencies.tasks SET title=title", + "DELETE FROM task_dependencies.tasks", + "UPDATE task_dependencies.projects SET account_id=account_id", + ]: + with psycopg.connect(autocommit=True) as connection: + with self.assertRaises(psycopg.Error): + connection.execute(sql) + + +if __name__ == "__main__": + unittest.main(verbosity=2) diff --git a/applications/task-dependencies/tests/requirements.txt b/applications/task-dependencies/tests/requirements.txt new file mode 100644 index 00000000..6c6cc80f --- /dev/null +++ b/applications/task-dependencies/tests/requirements.txt @@ -0,0 +1,3 @@ +psycopg==3.3.6 +psycopg-binary==3.3.6 +typing_extensions==4.16.0 diff --git a/applications/task-dependencies/tests/restart.py b/applications/task-dependencies/tests/restart.py new file mode 100644 index 00000000..5ebf4b47 --- /dev/null +++ b/applications/task-dependencies/tests/restart.py @@ -0,0 +1,119 @@ +"""Restart only this app's verified process, then compare its exact stored graph.""" + +import json +import os +from pathlib import Path +import signal +import subprocess +import time +import urllib.request +from acceptance import call, owner + +app = Path(__file__).resolve().parents[1] +evidence = Path(os.environ.get("EVIDENCE_DIR", "/tmp/task-dependencies-evidence")) +evidence.mkdir(parents=True, exist_ok=True) +project = "20000000-0000-4000-8000-000000000088" +with owner() as connection: + connection.execute( + "INSERT INTO task_dependencies.projects(id,account_id,name) VALUES (%s,%s,%s)", + (project, "10000000-0000-4000-8000-000000000001", "Persistence fixture"), + ) +tasks = [] +for title in ["Review draft", "Publish", "Notify team"]: + status, result = call( + f"/projects/{project}/tasks", method="POST", payload={"title": title} + ) + assert status == 201 + tasks.append(result) +for dependent, prerequisite in zip(tasks[1:], tasks): + status, _ = call( + f"/projects/{project}/tasks/{dependent['id']}/prerequisites", + method="POST", + payload={"prerequisite_id": prerequisite["id"]}, + ) + assert status == 200 +for task in tasks: + status, result = call( + f"/projects/{project}/tasks/{task['id']}/complete", method="POST", payload={} + ) + assert status == 200 and result["done"] and result["done_at"] +status, before = call(f"/projects/{project}") +assert status == 200 and len(before["tasks"]) == 3 and len(before["edges"]) == 2 +old = int(os.environ["API_PID"]) +proc = Path(f"/proc/{old}") +assert (proc / "cwd").resolve() == app +assert (proc / "exe").resolve() == app / "bin/task-api" +os.kill(old, signal.SIGTERM) +deadline = time.monotonic() + 20 +while proc.exists() and time.monotonic() < deadline: + time.sleep(0.1) +assert not proc.exists(), "Old process did not exit; refusing replacement." +print(f"Original PID {old} exited before replacement.", flush=True) +keys = [ + "PGHOST", + "PGPORT", + "PGDATABASE", + "PGUSER", + "PGPASSWORD", + "PGSSLMODE", + "PGSSLROOTCERT", + "PGCONNECT_TIMEOUT", + "TASKS_NORTH_TOKEN", + "TASKS_SOUTH_TOKEN", +] +environment = { + "HOME": os.environ["HOME"], + "PATH": "/usr/bin:/bin", + "LANG": "C.UTF-8", + **{key: os.environ[key] for key in keys}, +} +assert environment["PGUSER"] == "tasks_runtime" +child = subprocess.Popen( + [str(app / "bin/task-api")], + cwd=app, + env=environment, + stdout=(evidence / "replacement-private.log").open("ab"), + stderr=subprocess.STDOUT, + start_new_session=True, +) +(evidence / "replacement.pid").write_text(str(child.pid)) +deadline = time.monotonic() + 30 +while time.monotonic() < deadline: + assert child.poll() is None + try: + status, after = call(f"/projects/{project}") + if status == 200: + break + except OSError: + pass + time.sleep(0.2) +else: + raise AssertionError("Replacement not ready.") +assert before == after +status, completed = call( + f"/projects/{project}/tasks/{tasks[-1]['id']}/complete", method="POST", payload={} +) +assert status == 200 and completed == next( + task for task in before["tasks"] if task["id"] == tasks[-1]["id"] +) +child_keys = { + item.split(b"=", 1)[0].decode() + for item in Path(f"/proc/{child.pid}/environ").read_bytes().split(b"\0") + if item +} +assert not any("OWNER" in key or "TEST_" in key or "ADMIN" in key for key in child_keys) +(evidence / "persistence.json").write_text( + json.dumps( + { + "old_pid": old, + "new_pid": child.pid, + "before": before, + "after": after, + "idempotent_completion": completed, + }, + indent=2, + ) +) +print( + f"Replacement PID {child.pid}: exact project/revision, three done tasks/timestamps and two edges agree; repeat completion unchanged; runtime-only credentials." +)