diff --git a/.github/workflows/inventory-transfers.yml b/.github/workflows/inventory-transfers.yml new file mode 100644 index 00000000..03568099 --- /dev/null +++ b/.github/workflows/inventory-transfers.yml @@ -0,0 +1,30 @@ +name: Inventory Transfers checks +on: + pull_request: + paths: + - 'applications/inventory-transfers/**' + - '.github/workflows/inventory-transfers.yml' + push: + branches: [main] + paths: + - 'applications/inventory-transfers/**' + - '.github/workflows/inventory-transfers.yml' + workflow_dispatch: +permissions: + contents: read +jobs: + check: + runs-on: ubuntu-latest + defaults: + run: + working-directory: applications/inventory-transfers + steps: + - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 + - uses: dart-lang/setup-dart@6afc89df92d6eb3834022f73cd65adc8cdfcb92d # v1 + with: + sdk: '3.13.5' + - run: dart pub get --enforce-lockfile + - run: dart format --output=none --set-exit-if-changed bin lib test + - run: dart analyze + - run: dart test test/domain_test.dart test/http_test.dart + - run: dart compile exe bin/server.dart -o /tmp/inventory-transfers diff --git a/applications/inventory-transfers/.env.example b/applications/inventory-transfers/.env.example new file mode 100644 index 00000000..1ae3ddf7 --- /dev/null +++ b/applications/inventory-transfers/.env.example @@ -0,0 +1,12 @@ +PGHOST=your-cloud-postgres-host +PGPORT=5432 +PGDATABASE=postgres +PGUSER=transfers_app +PGPASSWORD=replace-me +PGSSLROOTCERT=/private/path/ca.pem +MIGRATION_PASSWORD=replace-with-separate-password +ADMIN_USER=cloud-created-user +ADMIN_PASSWORD=cloud-created-password +NORTH_TOKEN=replace-with-at-least-32-random-characters +SOUTH_TOKEN=replace-with-another-32-random-characters +PORT=4000 diff --git a/applications/inventory-transfers/.gitignore b/applications/inventory-transfers/.gitignore new file mode 100644 index 00000000..d886e48a --- /dev/null +++ b/applications/inventory-transfers/.gitignore @@ -0,0 +1,4 @@ +.dart_tool/ +.env +*.pem +build/ diff --git a/applications/inventory-transfers/README.md b/applications/inventory-transfers/README.md new file mode 100644 index 00000000..ec7f61b9 --- /dev/null +++ b/applications/inventory-transfers/README.md @@ -0,0 +1,165 @@ +# Inventory Transfers + +A small Dart and Shelf JSON API moves synthetic stock between warehouses within one organization. A retained request UUID makes a matching retry return the original immutable transfer. Balances are read separately because they may have changed since that transfer. + +This example uses ClickHouse Managed Postgres in ClickHouse Cloud. It demonstrates stock movement, not purchasing, shipping, payments or reconciliation with physical inventory. Two environment bearer tokens represent trusted organization operators. This is a loopback workbench, not a public identity system; a Flutter/mobile client must never contain database credentials or these shared tokens. + +Tested on 2 October 2026 with Dart 3.13.5 on Linux ARM64, Shelf 1.4.2, shelf_router 1.1.4, postgres 3.5.18 and PostgreSQL 18.6. `pubspec.lock` fixes transitive packages. + +## How a transfer commits + +`Transfers.submit()` runs every query on the `Pool.runTx` callback's transaction session: + +1. Insert the organization/request UUID with `ON CONFLICT DO NOTHING RETURNING`. +2. For a retained key, read the committed transfer in a fresh READ COMMITTED statement. Matching source/destination/SKU/units return it; changed payload returns 409. UUID case is normalized. +3. For a new key, lock both existing balance rows in warehouse order with `FOR NO KEY UPDATE`, check source stock and destination capacity, then debit and credit. The transfer and both writes commit together. + +The earlier insert's foreign keys acquire KEY SHARE locks on balances. Quantity-only `NO KEY UPDATE` locks are compatible with those checks while serializing quantity changes. Stronger `FOR UPDATE` upgrades can cause deadlocks between opposite-direction inserts. Both directions acquire the same pair in the same order. Exceptions escape the callback so the driver rolls back; the app never catches a uniqueness error and continues an aborted transaction. + +If a connection fails while returning a commit response, its outcome may be unknown. Retry the same retained request ID rather than inventing another one. This is a stored-record retry protocol, not a claim about exactly-once physical effects. + +## Create your Cloud fixture + +Cloud services incur charges until deleted. Use a dedicated service; the reset script destroys this example's schema/roles. See the [Managed Postgres quickstart](https://clickhouse.com/docs/products/managed-postgres/quickstart). + +```sh +umask 077 +# Use an existing private directory outside the checkout. +export ORG_ID=your-clickhouse-organization-id +clickhousectl cloud postgres create \ + --org-id "$ORG_ID" --name inventory-transfers-demo \ + --provider aws --region us-east-1 --size c6gd.large \ + --pg-version 18 --ha-type none --json > /private/path/create.json +``` + +Store the returned ID, hostname, username and one-time password privately. Check the current available size/region for your organization; this modest shape was supported by the tested service. Wait until `state` is `running`: + +```sh +export PG_ID=your-created-service-id +clickhousectl cloud postgres get "$PG_ID" --org-id "$ORG_ID" +clickhousectl cloud postgres certs get "$PG_ID" --org-id "$ORG_ID" \ + --output /private/path/ca.pem +``` + +Use `--output` to get a PEM file: the CLI automatically returns JSON to coding agents on stdout. The app creates a `SecurityContext(withTrustedRoots: false)` from this official CA and uses `SslMode.verifyFull`. In this driver `SslMode.require` ignores certificate errors and must not replace verification. The tested official bundle worked unchanged. Fetch a fresh certificate after Cloud certificate rotation. + +## Native setup + +Install stable Dart 3.13.5 for your Linux architecture and PostgreSQL's `psql` client. For the tested ARM64 VM: + +```sh +sudo apt-get update +sudo apt-get install -y curl ca-certificates unzip postgresql-client +curl -fsSLO https://storage.googleapis.com/dart-archive/channels/stable/release/3.13.5/sdk/dartsdk-linux-arm64-release.zip +curl -fsSLO https://storage.googleapis.com/dart-archive/channels/stable/release/3.13.5/sdk/dartsdk-linux-arm64-release.zip.sha256sum +sha256sum -c dartsdk-linux-arm64-release.zip.sha256sum +unzip -q dartsdk-linux-arm64-release.zip +export PATH="$PWD/dart-sdk/bin:$PATH" +cd /path/to/examples/applications/inventory-transfers +dart pub get --enforce-lockfile +``` + +Copy `.env.example` to a private file outside the checkout, fill in the Cloud fields and separate role passwords, and generate two distinct 32–128 character ASCII bearer tokens. Do not commit that file or certificates. `PGPASSWORD` holds the runtime password; `MIGRATION_PASSWORD` is separate. `ADMIN_USER`/`ADMIN_PASSWORD` come from the Cloud create response. Export the fields before calling child processes: + +```sh +set -a +source /private/path/setup.env +set +a +export PGSSLMODE=verify-full +APP_PASSWORD=$PGPASSWORD +export PGUSER=$ADMIN_USER PGPASSWORD=$ADMIN_PASSWORD +psql -X -v MIGRATION_PASSWORD="$MIGRATION_PASSWORD" -v APP_PASSWORD="$APP_PASSWORD" \ + -f sql/bootstrap.sql +export PGUSER=transfers_migration PGPASSWORD=$MIGRATION_PASSWORD +psql -X -f sql/migrate.sql +psql -X -f sql/grants.sql +psql -X -f sql/seed.sql +``` + +Bootstrap creates the schema and roles and revokes PUBLIC database CREATE/TEMP and public-schema CREATE on this dedicated fixture. The migration owner runs explicit versioned SQL under an advisory lock; migration and seed can be repeated. Seed inserts all warehouse/SKU combinations without resetting existing quantities. No startup migrations or automatic DDL run in the server. + +The runtime can SELECT balances/transfers, INSERT transfers, use the identity sequence, and UPDATE only balance quantities. It cannot edit/delete transfers, change balance keys, create tables/schema/temp tables, or read migration history. Composite foreign keys keep organization/SKU/endpoints together; CHECKs enforce distinct endpoints, positive bounded units and bounded nonnegative balances. The shared runtime credential remains trusted: direct quantity UPDATEs can bypass application conservation. The database is not an independently enforced ledger. + +## Run and try the API + +```sh +export PGUSER=transfers_app PGPASSWORD=$APP_PASSWORD +unset ADMIN_USER ADMIN_PASSWORD MIGRATION_PASSWORD +dart run bin/server.dart +``` + +The server binds `127.0.0.1:4000` by default; PORT can change that local port. It verifies connectivity at startup. SIGINT/SIGTERM stops the listener, finishes active work and closes the pool. Alternatively, build and run the native executable: + +```sh +mkdir -p build +dart compile exe bin/server.dart -o build/inventory-transfers +./build/inventory-transfers +``` + +Create `/private/path/runtime.env` with only `PGHOST`, `PGPORT`, `PGDATABASE`, `PGUSER=transfers_app`, its `PGPASSWORD`, `PGSSLROOTCERT`, `NORTH_TOKEN`, `SOUTH_TOKEN` and optional `PORT`. Copy the runtime values from your private setup file; exclude `ADMIN_USER`, `ADMIN_PASSWORD` and `MIGRATION_PASSWORD`. Keep both files mode 600 with `chmod 600 /private/path/setup.env /private/path/runtime.env`. In another runtime-only shell: + +```sh +set -a; source /private/path/runtime.env; set +a +curl -sS http://127.0.0.1:4000/balances \ + -H "Authorization: Bearer $NORTH_TOKEN" +curl -sS http://127.0.0.1:4000/transfers \ + -H "Authorization: Bearer $NORTH_TOKEN" -H 'Content-Type: application/json' \ + --data '{"request_id":"aaaaaaaa-1234-4567-890a-123456789abc","source":"depot","destination":"studio","sku":"bolts","units":25}' +``` + +Repeat that exact payload for 200/replay; change units under its key for 409. New transfers return 201. Unknown organization-scoped balances return 404; insufficient stock and destination capacity return 409. Tokens select the organization on the server; an organization body field is rejected. `/health` is public and reports process readiness after startup. + +| Route | Bounded behavior | +| --- | --- | +| GET `/balances?limit=25&after=depot%2Fbolts` | Warehouse/SKU ascending; tuple cursor matches that order | +| GET `/transfers?limit=25&before=123` | Numeric identity descending; IDs remain decimal strings | +| POST `/transfers` | Exactly five fields; UUID, ASCII keys and integer units 1–1,000,000 | + +Both lists default to 25 and cap at 100. `next_after`/`next_before` contain the last row's cursor; an additional empty page finishes traversal. Identity sequences can have gaps from retries/rollbacks, and their order is allocation order, not commit chronology. Retained rows are not automatically deleted. Lists are live reads, not a snapshot across pages. Transfer ordering qualifies the numeric table column: an unqualified `id` would select the `id::text` output alias and sort lexically. A targeted real API regression creates 12 transfers and pages across the 9/10/11/12 boundary. + +JSON bodies cap at 4 KiB, with a 3-second **inactivity** timeout between stream events. Authentication happens before reading the body; eight active requests are admitted and further requests return 503. Error responses close their HTTP connection to avoid reusing unread rejected body bytes. A client still streaming a rejected body may observe a disconnect rather than a JSON response; the acceptance helper sends an explicit Content-Length. The pool caps at 4 connections; postgres 3.5.18's 5-second connectTimeout covers pool wait plus remaining connection setup. Queries have 8-second client timeout, 5-second server statement timeout and 3-second lock timeout. Busy/transient or unknown-outcome errors use 503/retry_request_id. Tokens, passwords and query parameters are not logged. + +## Checks + +```sh +dart format --output=none --set-exit-if-changed bin lib test +dart analyze +dart test test/domain_test.dart test/http_test.dart +``` + +The three native checks cover numeric/identifier parsing and a real loopback HTTP preflight. For the destructive **dedicated Cloud fixture** acceptance, first reset using administrator credentials, then repeat the setup above. Restore private administrator/migration fields in the test shell; the application process launched by tests receives a runtime-only whitelist. + +```sh +set -a; source /private/path/setup.env; set +a +export APP_PASSWORD=$PGPASSWORD PGSSLMODE=verify-full +export PGUSER=$ADMIN_USER PGPASSWORD=$ADMIN_PASSWORD +psql -X -f sql/cleanup.sql +# Run bootstrap, migration twice, grants and seed twice as above. +export PGUSER=transfers_app PGPASSWORD=$APP_PASSWORD +openssl req -x509 -newkey rsa:2048 -nodes -days 2 -subj /CN=UnrelatedAcceptanceCA \ + -keyout /private/path/wrong-ca.key -out /private/path/wrong-ca.pem +export WRONG_CA=/private/path/wrong-ca.pem +export SERVER_EXECUTABLE=/private/path/transfers-server +export WORKER_EXECUTABLE=/private/path/transfers-worker +dart compile exe bin/server.dart -o "$SERVER_EXECUTABLE" +dart compile exe bin/worker.dart -o "$WORKER_EXECUTABLE" +RUN_CLOUD=true dart test test/cloud_test.dart --concurrency=1 --reporter=expanded +``` + +Cloud acceptance checks organization-scoped HTTP replay and invalid inputs, prefix-key pagination, native constraints/grant denials, owner-only credit-trigger rollback after debit, independent blocked workers, actual opposite-direction HTTP requests, competing withdrawals, capacity, certificate controls and original-process exit before durable restart/replay. Tests use owner-only balance resets between cases and controlled trigger/foreign-warehouse fixtures; these are test setup, not runtime permissions. TLS negatives assert the pinned driver's specific `BadCertificateException`: wrong CA and an IP target with mismatched certificate name receive the same real peer certificate, then the actual DNS endpoint verifies successfully. No verification bypass is used. + +## Cleanup + +Stop the app. Optional schema reset requires freshly restored administrator credentials, since runtime instructions unset them: + +```sh +set -a; source /private/path/setup.env; set +a +export PGUSER=$ADMIN_USER PGPASSWORD=$ADMIN_PASSWORD PGSSLMODE=verify-full +psql -X -f sql/cleanup.sql +clickhousectl cloud postgres delete "$PG_ID" --org-id "$ORG_ID" +clickhousectl cloud postgres list --org-id "$ORG_ID" +``` + +Confirm your exact ID is absent. Cleanup does not restore PUBLIC privileges or provide a production downgrade. + +Primary references: [Dart PostgreSQL client](https://pub.dev/packages/postgres), [pool transaction session](https://pub.dev/documentation/postgres/latest/postgres/Pool/runTx.html), [pool settings](https://pub.dev/documentation/postgres/latest/postgres/PoolSettings-class.html), [TLS modes](https://pub.dev/documentation/postgres/latest/postgres/SslMode.html), [Shelf](https://pub.dev/packages/shelf), [Dart native compilation](https://dart.dev/tools/dart-compile). diff --git a/applications/inventory-transfers/analysis_options.yaml b/applications/inventory-transfers/analysis_options.yaml new file mode 100644 index 00000000..b323461f --- /dev/null +++ b/applications/inventory-transfers/analysis_options.yaml @@ -0,0 +1,5 @@ +analyzer: + errors: + unused_import: error + unused_local_variable: error + dead_code: error diff --git a/applications/inventory-transfers/bin/server.dart b/applications/inventory-transfers/bin/server.dart new file mode 100644 index 00000000..db18c7fe --- /dev/null +++ b/applications/inventory-transfers/bin/server.dart @@ -0,0 +1,38 @@ +import 'dart:io'; + +import 'package:shelf/shelf_io.dart' as shelf_io; +import 'package:inventory_transfers/api.dart'; +import 'package:inventory_transfers/database.dart'; +import 'package:inventory_transfers/transfers.dart'; + +Future main() async { + if (requiredEnv('PGUSER') != 'transfers_app') + throw StateError('Runtime requires transfers_app'); + final port = int.parse(Platform.environment['PORT'] ?? '4000'); + if (port < 1024 || port > 65535) throw StateError('Invalid PORT'); + final pool = createPool(); + final api = Api(Transfers(pool), { + 'north': requiredEnv('NORTH_TOKEN'), + 'south': requiredEnv('SOUTH_TOKEN'), + }); + // Fail startup on an invalid database/TLS configuration, rather than on the first write. + await pool.execute('SELECT 1'); + final server = await shelf_io.serve( + api.handler, + InternetAddress.loopbackIPv4, + port, + ); + server.autoCompress = false; + stdout.writeln('Inventory Transfers listening on loopback port $port'); + var stopping = false; + Future stop() async { + if (stopping) return; + stopping = true; + await server.close(force: false); + await pool.close(); + exit(0); + } + + ProcessSignal.sigint.watch().listen((_) => stop()); + ProcessSignal.sigterm.watch().listen((_) => stop()); +} diff --git a/applications/inventory-transfers/bin/worker.dart b/applications/inventory-transfers/bin/worker.dart new file mode 100644 index 00000000..8959a163 --- /dev/null +++ b/applications/inventory-transfers/bin/worker.dart @@ -0,0 +1,26 @@ +import 'dart:convert'; +import 'dart:io'; + +import 'package:inventory_transfers/database.dart'; +import 'package:inventory_transfers/transfers.dart'; +import 'package:inventory_transfers/validation.dart'; + +// Independent acceptance worker: result is written atomically to a caller-chosen private file. +Future main(List arguments) async { + final input = jsonDecode(arguments[0]) as Map; + final pool = createPool(applicationName: arguments[2]); + Map output; + try { + final result = await Transfers(pool) + .submit(input.remove('org') as String, TransferInput.parse(input)); + output = {'transfer': result.transfer, 'replay': result.replay}; + } on ApiError catch (error) { + output = {'error': error.code}; + } catch (_) { + output = {'error': 'database_failure'}; + } + await pool.close(); + final temp = File('${arguments[1]}.tmp'); + await temp.writeAsString(jsonEncode(output)); + await temp.rename(arguments[1]); +} diff --git a/applications/inventory-transfers/lib/api.dart b/applications/inventory-transfers/lib/api.dart new file mode 100644 index 00000000..52b97195 --- /dev/null +++ b/applications/inventory-transfers/lib/api.dart @@ -0,0 +1,153 @@ +import 'dart:async'; +import 'dart:convert'; + +import 'package:postgres/postgres.dart'; +import 'package:shelf/shelf.dart'; +import 'package:shelf_router/shelf_router.dart'; + +import 'transfers.dart'; +import 'validation.dart'; + +Response jsonResponse(int status, Object body) => Response( + status, + body: jsonEncode(body), + headers: { + 'content-type': 'application/json; charset=utf-8', + 'cache-control': 'no-store', + // An early body rejection must not reuse a connection with unread request bytes. + if (status >= 400) 'connection': 'close', + }, +); + +bool equalToken(String a, String b) { + if (a.length != b.length) return false; + var difference = 0; + for (var i = 0; i < a.length; i++) { + difference |= a.codeUnitAt(i) ^ b.codeUnitAt(i); + } + return difference == 0; +} + +class Api { + final Transfers store; + final Map tokens; + var _active = 0; + Api(this.store, this.tokens) { + if (tokens.length != 2 || + tokens.values.toSet().length != 2 || + tokens.values.any( + (token) => !RegExp(r'^[!-~]{32,128}$').hasMatch(token), + )) { + throw StateError( + 'Configure two distinct 32–128 character organization tokens', + ); + } + } + + Handler get handler { + final router = Router(); + router.get('/health', (_) => jsonResponse(200, {'status': 'up'})); + router.get('/balances', (Request request) async { + final query = request.url.queryParameters; + if (query.keys.any((k) => !{'limit', 'after'}.contains(k))) + throw const ApiError(400, 'invalid_query'); + final after = query['after']; + if (after != null) { + final parts = after.split('/'); + if (parts.length != 2) throw const ApiError(400, 'invalid_cursor'); + key(parts[0]); + key(parts[1]); + } + final rows = await store.balances( + request.context['org'] as String, + limit(query['limit']), + after, + ); + return jsonResponse(200, { + 'rows': rows, + 'next_after': rows.isEmpty + ? null + : '${rows.last['warehouse']}/${rows.last['sku']}', + }); + }); + router.get('/transfers', (Request request) async { + final query = request.url.queryParameters; + if (query.keys.any((k) => !{'limit', 'before'}.contains(k))) + throw const ApiError(400, 'invalid_query'); + final rows = await store.list( + request.context['org'] as String, + limit(query['limit']), + cursor(query['before']), + ); + return jsonResponse(200, { + 'rows': rows, + 'next_before': rows.isEmpty ? null : rows.last['id'], + }); + }); + router.post('/transfers', (Request request) async { + final input = TransferInput.parse(await readJson(request)); + final result = await store.submit( + request.context['org'] as String, + input, + ); + return jsonResponse(result.replay ? 200 : 201, { + 'transfer': result.transfer, + 'replay': result.replay, + }); + }); + return (request) async { + if (request.url.path == 'health' && request.method == 'GET') + return router.call(request); + final authorization = request.headers['authorization'] ?? ''; + String? org; + for (final entry in tokens.entries) { + if (equalToken(authorization, 'Bearer ${entry.value}')) org = entry.key; + } + // Authentication and admission precede reading the body or scheduling database work. + if (org == null) return jsonResponse(401, {'error': 'unauthorized'}); + if (_active >= 8) return jsonResponse(503, {'error': 'busy'}); + _active++; + try { + return await router.call(request.change(context: {'org': org})); + } on ApiError catch (error) { + return jsonResponse(error.status, {'error': error.code}); + } on ServerException catch (error) { + if (error.code == '23503') + return jsonResponse(404, {'error': 'unknown_balance'}); + if (error.code == '55P03' || error.code == '57014') + return jsonResponse(503, {'error': 'retry_request_id'}); + return jsonResponse(500, {'error': 'database_failure'}); + } catch (_) { + // A lost commit response has an unknown outcome: retry the same retained request ID. + return jsonResponse(503, {'error': 'retry_request_id'}); + } finally { + _active--; + } + }; + } +} + +Future readJson(Request request) async { + if ((request.headers['content-type'] ?? '').split(';').first != + 'application/json') { + throw const ApiError(415, 'json_required'); + } + final length = int.tryParse(request.headers['content-length'] ?? ''); + if (length != null && length > 4096) + throw const ApiError(413, 'body_too_large'); + final bytes = []; + try { + await for (final chunk in request.read().timeout( + const Duration(seconds: 3), + )) { + if (bytes.length + chunk.length > 4096) + throw const ApiError(413, 'body_too_large'); + bytes.addAll(chunk); + } + return jsonDecode(utf8.decode(bytes)); + } on FormatException { + throw const ApiError(400, 'invalid_json'); + } on TimeoutException { + throw const ApiError(408, 'body_timeout'); + } +} diff --git a/applications/inventory-transfers/lib/database.dart b/applications/inventory-transfers/lib/database.dart new file mode 100644 index 00000000..c758c71b --- /dev/null +++ b/applications/inventory-transfers/lib/database.dart @@ -0,0 +1,42 @@ +import 'dart:io'; + +import 'package:postgres/postgres.dart'; + +String requiredEnv(String name, [Map? env]) { + final value = (env ?? Platform.environment)[name]; + if (value == null || value.isEmpty) throw StateError('Missing $name'); + return value; +} + +Endpoint endpoint({String? host, String? user, String? password}) => Endpoint( + host: host ?? requiredEnv('PGHOST'), + port: int.parse(Platform.environment['PGPORT'] ?? '5432'), + database: requiredEnv('PGDATABASE'), + username: user ?? requiredEnv('PGUSER'), + password: password ?? requiredEnv('PGPASSWORD'), +); + +SecurityContext trust(String path) => + SecurityContext(withTrustedRoots: false)..setTrustedCertificates(path); + +Pool createPool({String applicationName = 'inventory-transfers'}) => + Pool.withEndpoints( + [endpoint()], + settings: PoolSettings( + maxConnectionCount: 4, + maxConnectionAge: const Duration(minutes: 10), + connectTimeout: const Duration(seconds: 5), + queryTimeout: const Duration(seconds: 8), + sslMode: SslMode.verifyFull, + securityContext: trust(requiredEnv('PGSSLROOTCERT')), + applicationName: applicationName, + timeZone: 'UTC', + onOpen: (connection) async { + await connection.execute("SET statement_timeout = '5s'"); + await connection.execute("SET lock_timeout = '3s'"); + await connection.execute( + "SET idle_in_transaction_session_timeout = '10s'", + ); + }, + ), + ); diff --git a/applications/inventory-transfers/lib/transfers.dart b/applications/inventory-transfers/lib/transfers.dart new file mode 100644 index 00000000..3eaac929 --- /dev/null +++ b/applications/inventory-transfers/lib/transfers.dart @@ -0,0 +1,138 @@ +import 'package:postgres/postgres.dart'; + +import 'validation.dart'; + +const _projection = + 'id::text, request_id::text, source, destination, sku, units, created_at'; +Map transferRow(ResultRow row) => { + 'id': row[0], + 'request_id': row[1], + 'source': row[2], + 'destination': row[3], + 'sku': row[4], + 'units': row[5], + 'created_at': (row[6] as DateTime).toUtc().toIso8601String(), +}; + +class TransferResult { + final Map transfer; + final bool replay; + const TransferResult(this.transfer, this.replay); +} + +class Transfers { + final Pool pool; + const Transfers(this.pool); + + Future submit(String org, TransferInput input) => pool.runTx( + (tx) async { + // Every operation uses this transaction session, including the conflict read. + final inserted = await tx.execute( + Sql.named('''INSERT INTO transfers.transfers + (organization, request_id, source, destination, sku, units) + VALUES (@org, @request::uuid, @source, @destination, @sku, @units) + ON CONFLICT (organization, request_id) DO NOTHING + RETURNING $_projection'''), + parameters: input.parameters(org), + ); + if (inserted.isEmpty) { + // A fresh READ COMMITTED statement sees the row after the conflicting insert commits. + final existing = await tx.execute( + Sql.named('''SELECT $_projection FROM transfers.transfers + WHERE organization = @org AND request_id = @request::uuid'''), + parameters: {'org': org, 'request': input.requestId}, + ); + if (existing.isEmpty) throw StateError('Retained transfer missing'); + final row = transferRow(existing.single); + if (row['source'] != input.source || + row['destination'] != input.destination || + row['sku'] != input.sku || + row['units'] != input.units) { + throw const ApiError(409, 'request_conflict'); + } + return TransferResult(row, true); + } + // Both directions acquire the same pre-existing pair in warehouse-key order. + final balances = await tx.execute( + Sql.named( + '''SELECT warehouse, quantity + FROM transfers.balances WHERE organization = @org AND sku = @sku + AND warehouse IN (@source, @destination) ORDER BY warehouse FOR NO KEY UPDATE''', + ), + parameters: { + 'org': org, + 'sku': input.sku, + 'source': input.source, + 'destination': input.destination, + }, + ); + if (balances.length != 2) throw const ApiError(404, 'unknown_balance'); + final quantities = { + for (final row in balances) row[0] as String: row[1] as int, + }; + if (quantities[input.source]! < input.units) + throw const ApiError(409, 'insufficient_stock'); + if (quantities[input.destination]! > maxBalance - input.units) { + throw const ApiError(409, 'destination_capacity'); + } + await tx.execute( + Sql.named('''UPDATE transfers.balances SET quantity = quantity - @units + WHERE organization = @org AND sku = @sku AND warehouse = @source'''), + parameters: { + 'org': org, + 'sku': input.sku, + 'source': input.source, + 'units': input.units, + }, + ); + await tx.execute( + Sql.named('''UPDATE transfers.balances SET quantity = quantity + @units + WHERE organization = @org AND sku = @sku AND warehouse = @destination'''), + parameters: { + 'org': org, + 'sku': input.sku, + 'destination': input.destination, + 'units': input.units, + }, + ); + return TransferResult(transferRow(inserted.single), false); + }, + settings: TransactionSettings(isolationLevel: IsolationLevel.readCommitted), + ); + + Future>> balances( + String org, + int count, + String? after, + ) => pool.run((session) async { + final rows = await session.execute( + Sql.named('''SELECT warehouse, sku, quantity FROM transfers.balances + WHERE organization = @org AND (@afterWarehouse::text IS NULL OR + (warehouse, sku) > (@afterWarehouse::text, @afterSku::text)) + ORDER BY warehouse, sku LIMIT @count'''), + parameters: { + 'org': org, + 'afterWarehouse': after?.split('/')[0], + 'afterSku': after?.split('/')[1], + 'count': count, + }, + ); + return rows + .map((row) => {'warehouse': row[0], 'sku': row[1], 'quantity': row[2]}) + .toList(); + }); + + Future>> list( + String org, + int count, + String? before, + ) => pool.run((session) async { + final rows = await session.execute( + Sql.named('''SELECT $_projection FROM transfers.transfers + WHERE organization = @org AND (@before::bigint IS NULL OR id < @before::bigint) + ORDER BY transfers.transfers.id DESC LIMIT @count'''), + parameters: {'org': org, 'before': before, 'count': count}, + ); + return rows.map(transferRow).toList(); + }); +} diff --git a/applications/inventory-transfers/lib/validation.dart b/applications/inventory-transfers/lib/validation.dart new file mode 100644 index 00000000..61ab576e --- /dev/null +++ b/applications/inventory-transfers/lib/validation.dart @@ -0,0 +1,88 @@ +class ApiError implements Exception { + final int status; + final String code; + const ApiError(this.status, this.code); +} + +const maxBalance = 1000000000; +const maxUnits = 1000000; +final _key = RegExp(r'^[a-z][a-z0-9-]{0,31}$'); +final _uuid = RegExp( + r'^[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}$', +); + +String key(Object? value) { + if (value is! String || !_key.hasMatch(value)) { + throw const ApiError(400, 'invalid_identifier'); + } + return value; +} + +String uuid(Object? value) { + if (value is! String || value.length != 36 || !_uuid.hasMatch(value)) { + throw const ApiError(400, 'invalid_request_id'); + } + return value.toLowerCase(); +} + +class TransferInput { + final String requestId, source, destination, sku; + final int units; + const TransferInput( + this.requestId, + this.source, + this.destination, + this.sku, + this.units, + ); + + factory TransferInput.parse(Object? body) { + const fields = {'request_id', 'source', 'destination', 'sku', 'units'}; + if (body is! Map || + body.length != fields.length || + body.keys.any((field) => !fields.contains(field))) { + throw const ApiError(400, 'invalid_fields'); + } + final units = body['units']; + if (units is! int || units < 1 || units > maxUnits) { + throw const ApiError(400, 'invalid_units'); + } + final source = key(body['source']); + final destination = key(body['destination']); + if (source == destination) throw const ApiError(400, 'same_warehouse'); + return TransferInput( + uuid(body['request_id']), + source, + destination, + key(body['sku']), + units, + ); + } + + Map parameters(String organization) => { + 'org': organization, + 'request': requestId, + 'source': source, + 'destination': destination, + 'sku': sku, + 'units': units, + }; +} + +int limit(String? value) { + if (value == null) return 25; + if (!RegExp(r'^[0-9]{1,3}$').hasMatch(value)) + throw const ApiError(400, 'invalid_limit'); + final n = int.parse(value); + if (n < 1 || n > 100) throw const ApiError(400, 'invalid_limit'); + return n; +} + +String? cursor(String? value) { + if (value == null) return null; + if (!RegExp(r'^[1-9][0-9]{0,18}$').hasMatch(value) || + BigInt.parse(value) > BigInt.parse('9223372036854775807')) { + throw const ApiError(400, 'invalid_cursor'); + } + return value; +} diff --git a/applications/inventory-transfers/pubspec.lock b/applications/inventory-transfers/pubspec.lock new file mode 100644 index 00000000..5054c297 --- /dev/null +++ b/applications/inventory-transfers/pubspec.lock @@ -0,0 +1,421 @@ +# Generated by pub +# See https://dart.dev/tools/pub/glossary#lockfile +packages: + _fe_analyzer_shared: + dependency: transitive + description: + name: _fe_analyzer_shared + sha256: fdcd9f70f9eb80df3bc5ed0fa67280df2595fd481948df8a9ab082e6a40ad04b + url: "https://pub.dev" + source: hosted + version: "108.0.0" + analyzer: + dependency: transitive + description: + name: analyzer + sha256: a51c769bff3b6dfbe9d60199b8606d808290702a296bef0c26a4ca391d414e46 + url: "https://pub.dev" + source: hosted + version: "14.4.0" + args: + dependency: transitive + description: + name: args + sha256: d0481093c50b1da8910eb0bb301626d4d8eb7284aa739614d2b394ee09e3ea04 + url: "https://pub.dev" + source: hosted + version: "2.7.0" + async: + dependency: transitive + description: + name: async + sha256: e2eb0491ba5ddb6177742d2da23904574082139b07c1e33b8503b9f46f3e1a37 + url: "https://pub.dev" + source: hosted + version: "2.13.1" + boolean_selector: + dependency: transitive + description: + name: boolean_selector + sha256: "8aab1771e1243a5063b8b0ff68042d67334e3feab9e95b9490f9a6ebf73b42ea" + url: "https://pub.dev" + source: hosted + version: "2.1.2" + buffer: + dependency: transitive + description: + name: buffer + sha256: "389da2ec2c16283c8787e0adaede82b1842102f8c8aae2f49003a766c5c6b3d1" + url: "https://pub.dev" + source: hosted + version: "1.2.3" + charcode: + dependency: transitive + description: + name: charcode + sha256: fb0f1107cac15a5ea6ef0a6ef71a807b9e4267c713bb93e00e92d737cc8dbd8a + url: "https://pub.dev" + source: hosted + version: "1.4.0" + cli_config: + dependency: transitive + description: + name: cli_config + sha256: ac20a183a07002b700f0c25e61b7ee46b23c309d76ab7b7640a028f18e4d99ec + url: "https://pub.dev" + source: hosted + version: "0.2.0" + collection: + dependency: transitive + description: + name: collection + sha256: "2f5709ae4d3d59dd8f7cd309b4e023046b57d8a6c82130785d2b0e5868084e76" + url: "https://pub.dev" + source: hosted + version: "1.19.1" + convert: + dependency: transitive + description: + name: convert + sha256: b30acd5944035672bc15c6b7a8b47d773e41e2f17de064350988c5d02adb1c68 + url: "https://pub.dev" + source: hosted + version: "3.1.2" + coverage: + dependency: transitive + description: + name: coverage + sha256: "956a3de0725ca232ad353565a8290d3357592bf4250f6f298a185e2d949c5d3d" + url: "https://pub.dev" + source: hosted + version: "1.15.1" + crypto: + dependency: transitive + description: + name: crypto + sha256: c8ea0233063ba03258fbcf2ca4d6dadfefe14f02fab57702265467a19f27fadf + url: "https://pub.dev" + source: hosted + version: "3.0.7" + file: + dependency: transitive + description: + name: file + sha256: a3b4f84adafef897088c160faf7dfffb7696046cb13ae90b508c2cbc95d3b8d4 + url: "https://pub.dev" + source: hosted + version: "7.0.1" + frontend_server_client: + dependency: transitive + description: + name: frontend_server_client + sha256: "404dbf9bce29beafbe9dcb19219447b6f7d49900cd42b8d3b7af102b24508e25" + url: "https://pub.dev" + source: hosted + version: "4.1.0" + glob: + dependency: transitive + description: + name: glob + sha256: "218aeb56050c714f62a3182775320dfa04602b55074873e24e31bbd39bda96fb" + url: "https://pub.dev" + source: hosted + version: "2.2.0" + http_methods: + dependency: transitive + description: + name: http_methods + sha256: "6bccce8f1ec7b5d701e7921dca35e202d425b57e317ba1a37f2638590e29e566" + url: "https://pub.dev" + source: hosted + version: "1.1.1" + http_multi_server: + dependency: transitive + description: + name: http_multi_server + sha256: aa6199f908078bb1c5efb8d8638d4ae191aac11b311132c3ef48ce352fb52ef8 + url: "https://pub.dev" + source: hosted + version: "3.2.2" + http_parser: + dependency: transitive + description: + name: http_parser + sha256: "178d74305e7866013777bab2c3d8726205dc5a4dd935297175b19a23a2e66571" + url: "https://pub.dev" + source: hosted + version: "4.1.2" + io: + dependency: transitive + description: + name: io + sha256: "2635216ca6a737e60de577ffa1a48a0bec76ca8a62917cfc1bb88c14c570646f" + url: "https://pub.dev" + source: hosted + version: "1.1.0" + logging: + dependency: transitive + description: + name: logging + sha256: c8245ada5f1717ed44271ed1c26b8ce85ca3228fd2ffdb75468ab01979309d61 + url: "https://pub.dev" + source: hosted + version: "1.3.0" + matcher: + dependency: transitive + description: + name: matcher + sha256: "31bd099b47c10cd1aeb55146a2d46ce0277630ecef3f7dae54ad7873f36696cd" + url: "https://pub.dev" + source: hosted + version: "0.12.20" + meta: + dependency: transitive + description: + name: meta + sha256: "307249ce4ff29d58a18e97f6345f539382eb9c9c29ecda628900f31de0443dd9" + url: "https://pub.dev" + source: hosted + version: "1.19.0" + mime: + dependency: transitive + description: + name: mime + sha256: bd47de35f07e27267e69c8c8b22edf9473bfee170a60d60fcc93730c5144b7f6 + url: "https://pub.dev" + source: hosted + version: "2.1.0" + node_preamble: + dependency: transitive + description: + name: node_preamble + sha256: "6e7eac89047ab8a8d26cf16127b5ed26de65209847630400f9aefd7cd5c730db" + url: "https://pub.dev" + source: hosted + version: "2.0.2" + package_config: + dependency: transitive + description: + name: package_config + sha256: ffcf4cf3d6c0b74ac43708d9f56625506e8a68aa935abe9d267a7330f320eb5d + url: "https://pub.dev" + source: hosted + version: "3.0.0" + path: + dependency: transitive + description: + name: path + sha256: "75cca69d1490965be98c73ceaea117e8a04dd21217b37b292c9ddbec0d955bc5" + url: "https://pub.dev" + source: hosted + version: "1.9.1" + pool: + dependency: transitive + description: + name: pool + sha256: "4177f68c237ea2128d1bee66ac17b2ce05ba3dbaafcbdd54c5d40a39d0b6b11c" + url: "https://pub.dev" + source: hosted + version: "1.5.3" + postgres: + dependency: "direct main" + description: + name: postgres + sha256: dee84c58a69ea165c8a8daf7d755049f5868d4407062dde9305f914d2a33d0d3 + url: "https://pub.dev" + source: hosted + version: "3.5.18" + pub_semver: + dependency: transitive + description: + name: pub_semver + sha256: "261236774e8b1d69cfc6b9eabbc96c40f25e7a2d6b171f3385d4f65d5734fb24" + url: "https://pub.dev" + source: hosted + version: "2.2.1" + shelf: + dependency: "direct main" + description: + name: shelf + sha256: e7dd780a7ffb623c57850b33f43309312fc863fb6aa3d276a754bb299839ef12 + url: "https://pub.dev" + source: hosted + version: "1.4.2" + shelf_packages_handler: + dependency: transitive + description: + name: shelf_packages_handler + sha256: "89f967eca29607c933ba9571d838be31d67f53f6e4ee15147d5dc2934fee1b1e" + url: "https://pub.dev" + source: hosted + version: "3.0.2" + shelf_router: + dependency: "direct main" + description: + name: shelf_router + sha256: f5e5d492440a7fb165fe1e2e1a623f31f734d3370900070b2b1e0d0428d59864 + url: "https://pub.dev" + source: hosted + version: "1.1.4" + shelf_static: + dependency: transitive + description: + name: shelf_static + sha256: c87c3875f91262785dade62d135760c2c69cb217ac759485334c5857ad89f6e3 + url: "https://pub.dev" + source: hosted + version: "1.1.3" + shelf_web_socket: + dependency: transitive + description: + name: shelf_web_socket + sha256: "3632775c8e90d6c9712f883e633716432a27758216dfb61bd86a8321c0580925" + url: "https://pub.dev" + source: hosted + version: "3.0.0" + source_map_stack_trace: + dependency: transitive + description: + name: source_map_stack_trace + sha256: c0713a43e323c3302c2abe2a1cc89aa057a387101ebd280371d6a6c9fa68516b + url: "https://pub.dev" + source: hosted + version: "2.1.2" + source_maps: + dependency: transitive + description: + name: source_maps + sha256: "14c2945847669b44089bb1222f66873d7ff7103c58911917f2a63c5a62327898" + url: "https://pub.dev" + source: hosted + version: "0.10.14" + source_span: + dependency: transitive + description: + name: source_span + sha256: "56a02f1f4cd1a2d96303c0144c93bd6d909eea6bee6bf5a0e0b685edbd4c47ab" + url: "https://pub.dev" + source: hosted + version: "1.10.2" + stack_trace: + dependency: transitive + description: + name: stack_trace + sha256: "277654b3034d17ac6f9f1cb5595db011b1d5d41e8806866db28e0abaa101c490" + url: "https://pub.dev" + source: hosted + version: "1.12.2" + stream_channel: + dependency: transitive + description: + name: stream_channel + sha256: "969e04c80b8bcdf826f8f16579c7b14d780458bd97f56d107d3950fdbeef059d" + url: "https://pub.dev" + source: hosted + version: "2.1.4" + string_scanner: + dependency: transitive + description: + name: string_scanner + sha256: "921cd31725b72fe181906c6a94d987c78e3b98c2e205b397ea399d4054872b43" + url: "https://pub.dev" + source: hosted + version: "1.4.1" + term_glyph: + dependency: transitive + description: + name: term_glyph + sha256: "7f554798625ea768a7518313e58f83891c7f5024f88e46e7182a4558850a4b8e" + url: "https://pub.dev" + source: hosted + version: "1.2.2" + test: + dependency: "direct dev" + description: + name: test + sha256: de5d145b0afff7921e5e788a880f52d7e5f3ae24068a202f6fd3b58e4ba26323 + url: "https://pub.dev" + source: hosted + version: "1.32.0" + test_api: + dependency: transitive + description: + name: test_api + sha256: "0a10344e901e5b2e63819567951cb6a06673ed6b84f40462188ff5a0c41f371f" + url: "https://pub.dev" + source: hosted + version: "0.7.14" + test_core: + dependency: transitive + description: + name: test_core + sha256: "80f3fb49087454e07e7e07c67578cfdd156c8c3a5227d8b3f47c7b2d019c2e93" + url: "https://pub.dev" + source: hosted + version: "0.6.20" + typed_data: + dependency: transitive + description: + name: typed_data + sha256: f9049c039ebfeb4cf7a7104a675823cd72dba8297f264b6637062516699fa006 + url: "https://pub.dev" + source: hosted + version: "1.4.0" + vm_service: + dependency: transitive + description: + name: vm_service + sha256: "5f37239c4851efcef929cea7824e76df7f2f0970aef85d66bbc430afa40e72f0" + url: "https://pub.dev" + source: hosted + version: "15.3.0" + watcher: + dependency: transitive + description: + name: watcher + sha256: "1398c9f081a753f9226febe8900fce8f7d0a67163334e1c94a2438339d79d635" + url: "https://pub.dev" + source: hosted + version: "1.2.1" + web: + dependency: transitive + description: + name: web + sha256: "868d88a33d8a87b18ffc05f9f030ba328ffefba92d6c127917a2ba740f9cfe4a" + url: "https://pub.dev" + source: hosted + version: "1.1.1" + web_socket: + dependency: transitive + description: + name: web_socket + sha256: "34d64019aa8e36bf9842ac014bb5d2f5586ca73df5e4d9bf5c936975cae6982c" + url: "https://pub.dev" + source: hosted + version: "1.0.1" + web_socket_channel: + dependency: transitive + description: + name: web_socket_channel + sha256: d645757fb0f4773d602444000a8131ff5d48c9e47adfe9772652dd1a4f2d45c8 + url: "https://pub.dev" + source: hosted + version: "3.0.3" + webkit_inspection_protocol: + dependency: transitive + description: + name: webkit_inspection_protocol + sha256: "87d3f2333bb240704cd3f1c6b5b7acd8a10e7f0bc28c28dcf14e782014f4a572" + url: "https://pub.dev" + source: hosted + version: "1.2.1" + yaml: + dependency: transitive + description: + name: yaml + sha256: f67cdd8e07d3c6329146aaef1ba043542b3134c12489f553ca9a7435d1068aea + url: "https://pub.dev" + source: hosted + version: "3.1.4" +sdks: + dart: ">=3.13.5 <4.0.0" diff --git a/applications/inventory-transfers/pubspec.yaml b/applications/inventory-transfers/pubspec.yaml new file mode 100644 index 00000000..108567cb --- /dev/null +++ b/applications/inventory-transfers/pubspec.yaml @@ -0,0 +1,12 @@ +name: inventory_transfers +description: A bounded organization-scoped stock transfer API. +publish_to: none +version: 1.0.0 +environment: + sdk: '>=3.13.5 <4.0.0' +dependencies: + postgres: 3.5.18 + shelf: 1.4.2 + shelf_router: 1.1.4 +dev_dependencies: + test: 1.32.0 diff --git a/applications/inventory-transfers/sql/bootstrap.sql b/applications/inventory-transfers/sql/bootstrap.sql new file mode 100644 index 00000000..e5b314b5 --- /dev/null +++ b/applications/inventory-transfers/sql/bootstrap.sql @@ -0,0 +1,9 @@ +\set ON_ERROR_STOP on +SELECT format('CREATE ROLE transfers_migration LOGIN PASSWORD %L', :'MIGRATION_PASSWORD') +WHERE NOT EXISTS (SELECT FROM pg_roles WHERE rolname = 'transfers_migration') \gexec +SELECT format('CREATE ROLE transfers_app LOGIN PASSWORD %L', :'APP_PASSWORD') +WHERE NOT EXISTS (SELECT FROM pg_roles WHERE rolname = 'transfers_app') \gexec +CREATE SCHEMA IF NOT EXISTS transfers AUTHORIZATION transfers_migration; +REVOKE CREATE ON SCHEMA public FROM PUBLIC; +SELECT format('REVOKE CREATE, TEMP ON DATABASE %I FROM PUBLIC', current_database()) \gexec +SELECT format('GRANT CONNECT ON DATABASE %I TO transfers_migration, transfers_app', current_database()) \gexec diff --git a/applications/inventory-transfers/sql/cleanup.sql b/applications/inventory-transfers/sql/cleanup.sql new file mode 100644 index 00000000..da673cfd --- /dev/null +++ b/applications/inventory-transfers/sql/cleanup.sql @@ -0,0 +1,7 @@ +\set ON_ERROR_STOP on +DROP SCHEMA IF EXISTS transfers CASCADE; +SELECT format('DROP OWNED BY %I', rolname) FROM pg_roles +WHERE rolname IN ('transfers_app', 'transfers_migration') \gexec +DROP ROLE IF EXISTS transfers_app; +DROP ROLE IF EXISTS transfers_migration; +-- Dedicated-fixture reset only; PUBLIC privilege revocations remain in place. diff --git a/applications/inventory-transfers/sql/grants.sql b/applications/inventory-transfers/sql/grants.sql new file mode 100644 index 00000000..22e5853f --- /dev/null +++ b/applications/inventory-transfers/sql/grants.sql @@ -0,0 +1,6 @@ +\set ON_ERROR_STOP on +GRANT USAGE ON SCHEMA transfers TO transfers_app; +GRANT SELECT ON transfers.balances, transfers.transfers TO transfers_app; +GRANT UPDATE (quantity) ON transfers.balances TO transfers_app; +GRANT INSERT ON transfers.transfers TO transfers_app; +GRANT USAGE ON SEQUENCE transfers.transfers_id_seq TO transfers_app; diff --git a/applications/inventory-transfers/sql/migrate.sql b/applications/inventory-transfers/sql/migrate.sql new file mode 100644 index 00000000..f0351207 --- /dev/null +++ b/applications/inventory-transfers/sql/migrate.sql @@ -0,0 +1,54 @@ +\set ON_ERROR_STOP on +BEGIN; +SELECT pg_advisory_xact_lock(739188); +CREATE TABLE IF NOT EXISTS transfers.schema_versions ( + version integer PRIMARY KEY, + applied_at timestamptz NOT NULL DEFAULT now() +); +SELECT NOT EXISTS (SELECT FROM transfers.schema_versions WHERE version = 1) AS apply_v1 \gset +\if :apply_v1 +CREATE TABLE transfers.organizations ( + id text PRIMARY KEY CHECK (id ~ '^[a-z][a-z0-9-]{0,31}$'), + label text NOT NULL +); +CREATE TABLE transfers.warehouses ( + organization text NOT NULL REFERENCES transfers.organizations(id), + id text NOT NULL CHECK (id ~ '^[a-z][a-z0-9-]{0,31}$'), + label text NOT NULL, + PRIMARY KEY (organization, id) +); +CREATE TABLE transfers.skus ( + organization text NOT NULL REFERENCES transfers.organizations(id), + id text NOT NULL CHECK (id ~ '^[a-z][a-z0-9-]{0,31}$'), + label text NOT NULL, + PRIMARY KEY (organization, id) +); +CREATE TABLE transfers.balances ( + organization text NOT NULL, + warehouse text NOT NULL, + sku text NOT NULL, + quantity integer NOT NULL CHECK (quantity BETWEEN 0 AND 1000000000), + PRIMARY KEY (organization, warehouse, sku), + FOREIGN KEY (organization, warehouse) REFERENCES transfers.warehouses(organization, id), + FOREIGN KEY (organization, sku) REFERENCES transfers.skus(organization, id) +); +CREATE TABLE transfers.transfers ( + id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY, + organization text NOT NULL, + request_id uuid NOT NULL, + source text NOT NULL, + destination text NOT NULL, + sku text NOT NULL, + units integer NOT NULL CHECK (units BETWEEN 1 AND 1000000), + created_at timestamptz NOT NULL DEFAULT now(), + CONSTRAINT different_warehouses CHECK (source <> destination), + CONSTRAINT organization_request UNIQUE (organization, request_id), + CONSTRAINT source_balance FOREIGN KEY (organization, source, sku) + REFERENCES transfers.balances(organization, warehouse, sku), + CONSTRAINT destination_balance FOREIGN KEY (organization, destination, sku) + REFERENCES transfers.balances(organization, warehouse, sku) +); +CREATE INDEX transfers_org_id ON transfers.transfers (organization, id DESC); +INSERT INTO transfers.schema_versions (version) VALUES (1); +\endif +COMMIT; diff --git a/applications/inventory-transfers/sql/seed.sql b/applications/inventory-transfers/sql/seed.sql new file mode 100644 index 00000000..e1c0d3bf --- /dev/null +++ b/applications/inventory-transfers/sql/seed.sql @@ -0,0 +1,19 @@ +\set ON_ERROR_STOP on +BEGIN; +INSERT INTO transfers.organizations VALUES ('north', 'North workshop'), ('south', 'South workshop') +ON CONFLICT DO NOTHING; +INSERT INTO transfers.warehouses +SELECT org, warehouse, initcap(warehouse) FROM + (VALUES ('north'), ('south')) AS o(org), + (VALUES ('depot'), ('studio'), ('reserve')) AS w(warehouse) +ON CONFLICT DO NOTHING; +INSERT INTO transfers.skus +SELECT org, sku, initcap(sku) FROM + (VALUES ('north'), ('south')) AS o(org), + (VALUES ('bolts'), ('panels')) AS s(sku) +ON CONFLICT DO NOTHING; +INSERT INTO transfers.balances +SELECT w.organization, w.id, s.id, 1000 +FROM transfers.warehouses w JOIN transfers.skus s USING (organization) +ON CONFLICT DO NOTHING; +COMMIT; diff --git a/applications/inventory-transfers/test/cloud_test.dart b/applications/inventory-transfers/test/cloud_test.dart new file mode 100644 index 00000000..b1623fa0 --- /dev/null +++ b/applications/inventory-transfers/test/cloud_test.dart @@ -0,0 +1,527 @@ +import 'dart:async'; +import 'dart:convert'; +import 'dart:io'; + +import 'package:postgres/postgres.dart'; +import 'package:test/test.dart'; +import 'package:inventory_transfers/database.dart'; + +int sequence = 0; +String requestId() => + '10000000-0000-4000-8000-${(++sequence).toString().padLeft(12, '0')}'; +Map payload({ + String? id, + String source = 'depot', + String destination = 'studio', + int units = 10, +}) => { + 'request_id': id ?? requestId(), + 'source': source, + 'destination': destination, + 'sku': 'bolts', + 'units': units, +}; + +Future privileged(String user, String password) => Connection.open( + endpoint(user: user, password: password), + settings: ConnectionSettings( + sslMode: SslMode.verifyFull, + securityContext: trust(requiredEnv('PGSSLROOTCERT')), + connectTimeout: const Duration(seconds: 5), + queryTimeout: const Duration(seconds: 20), + applicationName: 'transfers-acceptance-control', + ), +); + +void main() { + group( + 'Cloud acceptance', + () { + late Connection owner, admin; + late Pool pool; + late Process server; + late String executable; + final client = HttpClient(); + final logs = []; + Future<(int, Map)> call( + String path, { + Map? body, + String? raw, + bool auth = true, + bool south = false, + }) async { + final request = await client.openUrl( + body != null || raw != null ? 'POST' : 'GET', + Uri.parse('http://127.0.0.1:4000$path'), + ); + if (auth) + request.headers.set( + 'authorization', + 'Bearer ${requiredEnv(south ? 'SOUTH_TOKEN' : 'NORTH_TOKEN')}', + ); + if (body != null || raw != null) { + request.headers.contentType = ContentType.json; + final bytes = utf8.encode(raw ?? jsonEncode(body)); + request.contentLength = bytes.length; + request.add(bytes); + } + print( + 'HTTP ${request.method} $path bodyBytes=${request.contentLength} auth=$auth', + ); + final response = await request.close(); + return ( + response.statusCode, + jsonDecode(await response.transform(utf8.decoder).join()) + as Map, + ); + } + + Future start() async { + // Whitelist the runtime environment; no administrator or migration credentials. + final environment = { + for (final name in [ + 'PGHOST', + 'PGPORT', + 'PGDATABASE', + 'PGUSER', + 'PGPASSWORD', + 'PGSSLROOTCERT', + 'NORTH_TOKEN', + 'SOUTH_TOKEN', + ]) + name: requiredEnv(name), + 'PORT': '4000', + }; + server = await Process.start( + executable, + [], + environment: environment, + includeParentEnvironment: false, + ); + server.stdout.transform(utf8.decoder).listen(logs.add); + server.stderr.transform(utf8.decoder).listen(logs.add); + for (var attempt = 0; attempt < 100; attempt++) { + try { + final response = await call('/health', auth: false); + if (response.$1 == 200) return; + } catch (_) {} + await Future.delayed(const Duration(milliseconds: 100)); + } + throw StateError('Server did not become ready'); + } + + Future> quantities() async => (await owner.execute( + "SELECT quantity FROM transfers.balances WHERE organization='north' AND sku='bolts' AND warehouse IN ('depot','studio') ORDER BY warehouse", + )).map((row) => row[0] as int).toList(); + Future sql(String statement) async { + await owner.execute(statement); + } + + Future hasTransfer(String id) async => (await owner.execute( + Sql.named( + "SELECT 1 FROM transfers.transfers WHERE organization='north' AND request_id=@id::uuid", + ), + parameters: {'id': id}, + )).isNotEmpty; + Future denied(String statement, String code) async { + await expectLater( + pool.execute(statement), + throwsA( + isA().having((e) => e.code, 'SQLSTATE', code), + ), + ); + } + + Future worker( + Map body, + String name, + String output, + ) => Process.start(requiredEnv('WORKER_EXECUTABLE'), [ + jsonEncode({...body, 'org': 'north'}), + output, + name, + ]); + Future> finish(Process process, String path) async { + expect(await process.exitCode.timeout(const Duration(seconds: 15)), 0); + return jsonDecode(await File(path).readAsString()) + as Map; + } + + Future blocked(List names) async { + for (var attempt = 0; attempt < 100; attempt++) { + await admin.execute('SELECT pg_stat_clear_snapshot()'); + final rows = await admin.execute( + Sql.named( + '''SELECT application_name FROM pg_stat_activity + WHERE application_name=ANY(@names::text[]) AND cardinality(pg_blocking_pids(pid))>0''', + ), + parameters: {'names': names}, + ); + if (rows.length == names.length) return; + await Future.delayed(const Duration(milliseconds: 20)); + } + fail('Independent workers were not observed blocked'); + } + + setUpAll(() async { + executable = requiredEnv('SERVER_EXECUTABLE'); + owner = await privileged( + 'transfers_migration', + requiredEnv('MIGRATION_PASSWORD'), + ); + admin = await privileged( + requiredEnv('ADMIN_USER'), + requiredEnv('ADMIN_PASSWORD'), + ); + pool = createPool(applicationName: 'transfers-acceptance'); + print( + 'Actual PostgreSQL: ${(await owner.execute('SHOW server_version')).single.single}', + ); + await start(); + }); + setUp(() async { + await sql("UPDATE transfers.balances SET quantity=1000"); + }); + tearDownAll(() async { + server.kill(ProcessSignal.sigterm); + expect(await server.exitCode.timeout(const Duration(seconds: 15)), 0); + print('Server output: ${logs.join()}'); + client.close(force: true); + await pool.close(); + await owner.close(); + await admin.close(); + }); + + test('HTTP organization scope, retained replay, malformed input and tuple pagination', () async { + final body = payload(); + final created = await call('/transfers', body: body); + expect(created.$1, 201); + expect(created.$2['replay'], false); + final repeated = await call( + '/transfers', + body: { + ...body, + 'request_id': (body['request_id'] as String).toUpperCase(), + }, + ); + expect(repeated.$1, 200); + expect(repeated.$2['transfer'], created.$2['transfer']); + expect( + (await call('/transfers', body: {...body, 'units': 11})).$1, + 409, + ); + expect((await call('/transfers', body: body, south: true)).$1, 201); + expect( + (await call( + '/transfers', + body: {...body, 'organization': 'south'}, + )).$1, + 400, + ); + expect( + (await call( + '/transfers', + body: {...payload(), 'source': 'foreign-only'}, + )).$1, + 404, + ); + for (final value in [true, 1.5, 0, 1000001, '10']) { + expect( + (await call('/transfers', body: {...payload(), 'units': value})).$1, + 400, + ); + } + expect((await call('/transfers', raw: 'x' * 5000)).$1, 413); + expect( + (await call('/transfers', raw: 'x' * 5000, auth: false)).$1, + 401, + ); + expect((await call('/transfers', raw: '{bad')).$1, 400); + expect( + (await call( + '/transfers', + body: {...payload(), 'sku': 'bolts\u0000'}, + )).$1, + 400, + ); + await sql( + "INSERT INTO transfers.warehouses VALUES ('north','a','A'),('north','a-branch','Branch'),('south','foreign-only','Foreign')", + ); + await sql( + "INSERT INTO transfers.balances SELECT w.organization,w.id,s.id,1000 FROM transfers.warehouses w JOIN transfers.skus s USING(organization) WHERE w.id IN ('a','a-branch','foreign-only')", + ); + final expected = (await owner.execute( + "SELECT warehouse || '/' || sku FROM transfers.balances WHERE organization='north' ORDER BY warehouse,sku", + )).map((r) => r[0]).toList(); + final seen = []; + String? after; + for (var page = 0; page < 30; page++) { + final result = await call( + '/balances?limit=1${after == null ? '' : '&after=${Uri.encodeQueryComponent(after)}'}', + ); + expect(result.$1, 200); + final rows = result.$2['rows'] as List; + if (rows.isEmpty) break; + final row = rows.single as Map; + seen.add('${row['warehouse']}/${row['sku']}'); + after = result.$2['next_after'] as String; + } + expect(seen, expected); + expect(seen.toSet().length, seen.length); + expect((await call('/balances?limit=101')).$1, 400); + expect((await call('/transfers?before=9223372036854775808')).$1, 400); + final north = (await call('/transfers')).$2['rows'] as List; + final south = + (await call('/transfers', south: true)).$2['rows'] as List; + expect(north.length, 1); + expect(south.length, 1); + expect(north.single['id'], isNot(south.single['id'])); + // Now the foreign warehouse really exists, but only in the other organization. + expect( + (await call( + '/transfers', + body: {...payload(), 'source': 'foreign-only'}, + )).$1, + 404, + ); + }); + + test('numeric identity pagination across two-digit IDs', () async { + final created = []; + for (var n = 0; n < 12; n++) { + final result = await call('/transfers', body: payload(units: 1)); + expect(result.$1, 201); + created.add((result.$2['transfer'] as Map)['id'] as String); + } + print('Created numeric pagination IDs: ${created.join(',')}'); + expect(created.toSet().length, 12); + // Original unqualified alias sorts text: retain a real PostgreSQL control. + final oldOrder = (await owner.execute( + "SELECT id::text AS id FROM generate_series(1,12) AS fixture(id) ORDER BY id DESC", + )).map((row) => row[0] as String).toList(); + expect(oldOrder, isNot(List.generate(12, (n) => '${12 - n}'))); + final expected = (await owner.execute( + "SELECT id::text FROM transfers.transfers WHERE organization='north' ORDER BY transfers.transfers.id DESC", + )).map((row) => row[0] as String).toList(); + final seen = []; + String? before; + for (var page = 0; page < 8; page++) { + final result = await call( + '/transfers?limit=2${before == null ? '' : '&before=$before'}', + ); + expect(result.$1, 200); + final rows = result.$2['rows'] as List; + if (rows.isEmpty) break; + seen.addAll(rows.map((row) => (row as Map)['id'] as String)); + before = result.$2['next_before'] as String; + } + expect(seen, expected); + expect(seen.toSet().length, expected.length); + expect(await quantities(), [988, 1012]); + }); + + test('native constraints and restricted runtime grants', () async { + await denied("UPDATE transfers.transfers SET units=1", '42501'); + await denied("DELETE FROM transfers.transfers", '42501'); + await denied("UPDATE transfers.balances SET warehouse='x'", '42501'); + await denied("CREATE TABLE transfers.forbidden (id int)", '42501'); + await denied("CREATE SCHEMA forbidden", '42501'); + await denied("CREATE TEMP TABLE forbidden (id int)", '42501'); + await denied("SELECT * FROM transfers.schema_versions", '42501'); + await denied("UPDATE transfers.balances SET quantity=-1", '23514'); + await denied( + "INSERT INTO transfers.transfers (organization,request_id,source,destination,sku,units) VALUES ('north','20000000-0000-4000-8000-000000000001','depot','foreign-only','bolts',1)", + '23503', + ); + await denied( + "INSERT INTO transfers.transfers (organization,request_id,source,destination,sku,units) VALUES ('north','20000000-0000-4000-8000-000000000002','depot','depot','bolts',1)", + '23514', + ); + expect(await quantities(), [1000, 1000]); + }); + + test( + 'failure after debit rolls back both balances and retained transfer', + () async { + final body = payload(); + final before = await quantities(); + await sql( + r'''CREATE FUNCTION transfers.reject_credit() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN + IF NEW.organization='north' AND NEW.warehouse='studio' AND NEW.sku='bolts' AND NEW.quantity>OLD.quantity THEN + RAISE EXCEPTION 'acceptance credit failure'; END IF; RETURN NEW; END $$''', + ); + await sql( + 'CREATE TRIGGER acceptance_credit BEFORE UPDATE ON transfers.balances FOR EACH ROW EXECUTE FUNCTION transfers.reject_credit()', + ); + try { + expect((await call('/transfers', body: body)).$1, 500); + expect(await quantities(), before); + expect(await hasTransfer(body['request_id'] as String), false); + } finally { + await sql('DROP TRIGGER acceptance_credit ON transfers.balances'); + await sql('DROP FUNCTION transfers.reject_credit()'); + } + expect((await call('/transfers', body: body)).$1, 201); + expect(await quantities(), [990, 1010]); + }, + ); + + test( + 'independent same-key workers block and return one retained record', + () async { + final body = payload(); + final directory = await Directory.systemTemp.createTemp( + 'transfers-same-', + ); + await admin.execute('BEGIN'); + await admin.execute( + "SELECT 1 FROM transfers.balances WHERE organization='north' AND sku='bolts' AND warehouse IN ('depot','studio') ORDER BY warehouse FOR NO KEY UPDATE", + ); + final first = await worker( + body, + 'transfers-same-a', + '${directory.path}/a.json', + ); + final second = await worker( + body, + 'transfers-same-b', + '${directory.path}/b.json', + ); + try { + await blocked(['transfers-same-a', 'transfers-same-b']); + } finally { + await admin.execute('COMMIT'); + } + final a = await finish(first, '${directory.path}/a.json'); + final b = await finish(second, '${directory.path}/b.json'); + expect(a['transfer'], b['transfer']); + expect([a['replay'], b['replay']], containsAll([true, false])); + expect(await quantities(), [990, 1010]); + await directory.delete(recursive: true); + }, + ); + + test('independent opposite-direction contention and actual HTTP concurrency conserve sum', () async { + final directory = await Directory.systemTemp.createTemp( + 'transfers-opposite-', + ); + await admin.execute('BEGIN'); + await admin.execute( + "SELECT 1 FROM transfers.balances WHERE organization='north' AND sku='bolts' AND warehouse IN ('depot','studio') ORDER BY warehouse FOR NO KEY UPDATE", + ); + final a = await worker( + payload(units: 17), + 'transfers-opposite-a', + '${directory.path}/a.json', + ); + final b = await worker( + payload(source: 'studio', destination: 'depot', units: 23), + 'transfers-opposite-b', + '${directory.path}/b.json', + ); + try { + await blocked(['transfers-opposite-a', 'transfers-opposite-b']); + } finally { + await admin.execute('COMMIT'); + } + expect((await finish(a, '${directory.path}/a.json'))['error'], isNull); + expect((await finish(b, '${directory.path}/b.json'))['error'], isNull); + expect(await quantities(), [1006, 994]); + final results = await Future.wait([ + call('/transfers', body: payload(units: 11)), + call( + '/transfers', + body: payload(source: 'studio', destination: 'depot', units: 7), + ), + ]); + expect(results.map((r) => r.$1), [201, 201]); + expect(await quantities(), [1002, 998]); + expect((await quantities()).reduce((a, b) => a + b), 2000); + await directory.delete(recursive: true); + }); + + test('competing withdrawals cannot overspend and destination overflow leaves no record', () async { + await sql( + "UPDATE transfers.balances SET quantity=10 WHERE organization='north' AND sku='bolts' AND warehouse='depot'", + ); + final results = await Future.wait([ + call('/transfers', body: payload(units: 7)), + call('/transfers', body: payload(units: 7)), + ]); + expect(results.map((r) => r.$1).toList()..sort(), [201, 409]); + expect(await quantities(), [3, 1007]); + await sql( + "UPDATE transfers.balances SET quantity=1000000000 WHERE organization='north' AND sku='bolts' AND warehouse='studio'", + ); + final body = payload(units: 1); + expect( + (await call('/transfers', body: body)).$2['error'], + 'destination_capacity', + ); + expect(await quantities(), [3, 1000000000]); + expect(await hasTransfer(body['request_id'] as String), false); + }); + + test('real driver certificate and hostname failures with positive same-endpoint control', () async { + Future> reject(Endpoint target, String ca) async { + try { + final connection = await Connection.open( + target, + settings: ConnectionSettings( + sslMode: SslMode.verifyFull, + securityContext: trust(ca), + connectTimeout: const Duration(seconds: 5), + ), + ); + await connection.close(); + } on BadCertificateException catch (error) { + // The pinned driver reports verification through this specific exception. + return error.certificate.sha1; + } + fail('Invalid TLS unexpectedly succeeded'); + } + + final wrongCaCertificate = await reject( + endpoint(), + requiredEnv('WRONG_CA'), + ); + final address = (await InternetAddress.lookup(requiredEnv('PGHOST'))) + .firstWhere((a) => a.type == InternetAddressType.IPv4) + .address; + final wrongNameCertificate = await reject( + endpoint(host: address), + requiredEnv('PGSSLROOTCERT'), + ); + expect(wrongNameCertificate, wrongCaCertificate); + final connection = await Connection.open( + endpoint(), + settings: ConnectionSettings( + sslMode: SslMode.verifyFull, + securityContext: trust(requiredEnv('PGSSLROOTCERT')), + ), + ); + expect((await connection.execute('SELECT 1')).single.single, 1); + await connection.close(); + }); + + test('original compiled process exits before replacement and persisted request replays', () async { + final body = payload(units: 31); + final original = await call('/transfers', body: body); + expect(original.$1, 201); + final before = await quantities(); + final oldPid = server.pid; + server.kill(ProcessSignal.sigterm); + expect(await server.exitCode.timeout(const Duration(seconds: 15)), 0); + await start(); + expect(server.pid, isNot(oldPid)); + final replay = await call('/transfers', body: body); + expect(replay.$1, 200); + expect(replay.$2['transfer'], original.$2['transfer']); + expect(await quantities(), before); + expect((await call('/balances')).$1, 200); + }); + }, + skip: Platform.environment['RUN_CLOUD'] != 'true', + timeout: const Timeout(Duration(minutes: 3)), + ); +} diff --git a/applications/inventory-transfers/test/domain_test.dart b/applications/inventory-transfers/test/domain_test.dart new file mode 100644 index 00000000..7edcac27 --- /dev/null +++ b/applications/inventory-transfers/test/domain_test.dart @@ -0,0 +1,57 @@ +import 'dart:convert'; + +import 'package:test/test.dart'; +import 'package:inventory_transfers/validation.dart'; + +Map payload() => { + 'request_id': 'ABCDEF00-1234-4567-890A-123456789ABC', + 'source': 'depot', + 'destination': 'studio', + 'sku': 'bolts', + 'units': 10, +}; +void main() { + test('UUID normalization and exact bounded integers', () { + final value = TransferInput.parse(payload()); + expect(value.requestId, 'abcdef00-1234-4567-890a-123456789abc'); + for (final units in [true, 1.5, 0, -1, maxUnits + 1, '10', null]) { + expect( + () => TransferInput.parse({...payload(), 'units': units}), + throwsA(isA()), + ); + } + expect( + TransferInput.parse({...payload(), 'units': maxUnits}).units, + maxUnits, + ); + expect( + () => TransferInput.parse( + jsonDecode( + '{"request_id":"abcdef00-1234-4567-890a-123456789abc","source":"depot","destination":"studio","sku":"bolts","units":1.0}', + ), + ), + throwsA(isA()), + ); + }); + test( + 'strict fields, identifiers, distinct endpoints and bounded cursors', + () { + for (final body in [ + null, + [], + {...payload(), 'organization': 'south'}, + {...payload(), 'source': 'studio'}, + {...payload(), 'sku': 'bolts\u0000'}, + {...payload(), 'request_id': 'x' * 5000}, + ]) { + expect(() => TransferInput.parse(body), throwsA(isA())); + } + expect(cursor('9223372036854775807'), '9223372036854775807'); + for (final c in ['0', '-1', '9223372036854775808', '1' * 1000]) { + expect(() => cursor(c), throwsA(isA())); + } + expect(limit('100'), 100); + expect(() => limit('101'), throwsA(isA())); + }, + ); +} diff --git a/applications/inventory-transfers/test/http_test.dart b/applications/inventory-transfers/test/http_test.dart new file mode 100644 index 00000000..61015afc --- /dev/null +++ b/applications/inventory-transfers/test/http_test.dart @@ -0,0 +1,58 @@ +import 'dart:convert'; +import 'dart:io'; + +import 'package:postgres/postgres.dart'; +import 'package:shelf/shelf_io.dart' as shelf_io; +import 'package:test/test.dart'; +import 'package:inventory_transfers/api.dart'; +import 'package:inventory_transfers/transfers.dart'; + +void main() { + test( + 'real HTTP health, authentication before body and strict JSON boundaries', + () async { + // Lazy pool: these preflight routes reject before any database operation. + final pool = Pool.withEndpoints([ + Endpoint(host: '127.0.0.1', database: 'unused'), + ]); + final token = 'n' * 32; + final api = Api(Transfers(pool), {'north': token, 'south': 's' * 32}); + final server = await shelf_io.serve( + api.handler, + InternetAddress.loopbackIPv4, + 0, + ); + final client = HttpClient(); + try { + Future call(String path, {String? body, bool auth = true}) async { + final request = await client.openUrl( + body == null ? 'GET' : 'POST', + Uri.parse('http://127.0.0.1:${server.port}$path'), + ); + if (auth) request.headers.set('authorization', 'Bearer $token'); + if (body != null) { + request.headers.contentType = ContentType.json; + request.write(body); + } + final response = await request.close(); + await response.drain(); + return response.statusCode; + } + + expect(await call('/health', auth: false), 200); + expect(await call('/transfers', auth: false, body: 'x' * 5000), 401); + expect(await call('/transfers', body: 'x' * 5000), 413); + expect(await call('/transfers', body: '{bad'), 400); + expect( + await call('/transfers', body: jsonEncode({'units': true})), + 400, + ); + expect(await call('/balances?limit=101'), 400); + } finally { + client.close(force: true); + await server.close(force: true); + await pool.close(); + } + }, + ); +}