diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..b78baf9 --- /dev/null +++ b/.env.example @@ -0,0 +1,57 @@ +# Hymical Forms configuration. +# +# Copy this file to `.env`, or set the same variables in your process environment. +# FORMS_DATABASE_URL is required. Everything below it is optional and shown with +# its built-in default; uncomment the lines you want to change. + +# SQLAlchemy database URL. PostgreSQL is the intended production database. +FORMS_DATABASE_URL=postgresql+psycopg://forms:forms@localhost:5432/forms + +# SQLite is supported for local experimentation and backs the test suite. It is +# not a supported production target. +# FORMS_DATABASE_URL=sqlite:///./forms.db + +# Largest request body accepted, in bytes. File uploads are not supported, so +# this only needs to accommodate text form fields. +# FORMS_MAX_BODY_BYTES=262144 + +# Largest number of name/value pairs accepted in one submission. A repeated +# field name (a checkbox group) counts once per submitted value. +# FORMS_MAX_FIELDS=100 + +# Largest field name accepted, in characters. +# FORMS_MAX_FIELD_NAME_LENGTH=128 + +# Largest field value accepted, in characters. +# FORMS_MAX_FIELD_VALUE_LENGTH=16384 + +# How long to wait for a webhook destination to accept a connection, in seconds. +# FORMS_WEBHOOK_CONNECT_TIMEOUT_SECONDS=5 + +# How long to wait for a webhook destination to respond, in seconds. +# FORMS_WEBHOOK_READ_TIMEOUT_SECONDS=10 + +# How many delivery attempts a submission gets before it is given up on. +# FORMS_WEBHOOK_MAX_ATTEMPTS=5 + +# Wait before the second attempt, in seconds. Each later wait doubles it. +# FORMS_WEBHOOK_RETRY_INITIAL_SECONDS=10 + +# Cap on the wait between attempts, however far the backoff has doubled. +# FORMS_WEBHOOK_RETRY_MAX_SECONDS=3600 + +# How long the worker waits before looking for due deliveries again, in seconds. +# FORMS_WORKER_POLL_SECONDS=1 + +# How many deliveries a worker claims at once. +# FORMS_WORKER_BATCH_SIZE=10 + +# How long a worker's claim on a delivery holds, in seconds. After this the +# delivery becomes claimable again, which is how work is recovered from a worker +# that died holding it. +# FORMS_WORKER_LEASE_SECONDS=60 + +# Permit webhook destinations on loopback and private addresses, so a webhook can +# point at a server on your own machine. Development only: enabling this in +# production lets anyone who can create an endpoint reach your internal network. +# FORMS_ALLOW_PRIVATE_WEBHOOK_TARGETS=false diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml new file mode 100644 index 0000000..63161a0 --- /dev/null +++ b/.github/workflows/ci.yml @@ -0,0 +1,32 @@ +name: CI + +on: + push: + branches: [main] + pull_request: + +permissions: + contents: read + +concurrency: + group: ci-${{ github.ref }} + cancel-in-progress: true + +jobs: + test: + runs-on: ubuntu-latest + strategy: + fail-fast: false + matrix: + python-version: ["3.11", "3.12", "3.13"] + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-python@v5 + with: + python-version: ${{ matrix.python-version }} + cache: pip + - run: pip install -e ".[dev]" + - run: ruff check . + - run: ruff format --check . + - run: mypy + - run: pytest diff --git a/.gitignore b/.gitignore index 83972fa..d3323fe 100644 --- a/.gitignore +++ b/.gitignore @@ -61,6 +61,12 @@ local_settings.py db.sqlite3 db.sqlite3-journal +# Local SQLite databases, e.g. FORMS_DATABASE_URL=sqlite:///./forms.db +*.db +*.db-journal +*.sqlite +*.sqlite3 + # Flask stuff: instance/ .webassets-cache diff --git a/README.md b/README.md index f9cc2d9..760ecdb 100644 --- a/README.md +++ b/README.md @@ -1,2 +1,647 @@ -# forms -Reliable form submission infrastructure with validation, storage, webhooks, retries, and delivery tracking +
+
+
+ Reliable form ingestion and webhook delivery for developers. +
+ +## The problem + +Every project with a contact form, a waitlist, or a feedback box ends up needing +the same small backend: something that accepts an HTML form POST, validates it, +stores it, and forwards it somewhere useful. Writing that once is easy; running +it reliably, with retries, delivery logs, spam handling and retention rules, +is not. Hymical Forms is intended to be that backend, self-hostable and +open-source. + +## Project status + +**Early development.** This build registers endpoints, stores the submissions +sent to them together with the durable obligation to deliver them, and runs a +separate worker that performs the signed webhook delivery and retries it. + +Endpoint and webhook configuration is completely unauthenticated: anyone who can +reach the API can create an endpoint pointing anywhere. There is no rate limiting +and no spam protection, so do not expose this to the public internet. + +| Capability | Status | +| ----------------------------- | ------------------------- | +| Health endpoint | Implemented | +| Form ingestion + validation | Implemented | +| Request limits + error model | Implemented | +| Endpoint registry | Implemented | +| Submission persistence | Implemented | +| Idempotent retries | Implemented | +| Signed webhook delivery | Implemented | +| Durable delivery queue | Implemented | +| Retries with backoff | Implemented | +| API keys / authentication | **Not implemented** | +| Manual delivery replay | **Not implemented** | +| Rate limiting, spam handling | **Not implemented** | +| Schema migrations | **Not implemented** | +| Export, retention, dashboards | **Not implemented** | + +## Requirements + +- Python 3.11 or newer +- PostgreSQL, which is the intended production database + +SQLite is supported for local experimentation and backs the test suite. It is +not a supported production target. + +## Install + +```bash +python -m venv .venv && . .venv/bin/activate && pip install -e ".[dev]" +``` + +On Windows, activate with `.venv\Scripts\activate` instead. + +## Configure + +`FORMS_DATABASE_URL` is required and has no default. Set it in the environment +or in a `.env` file in the working directory: + +```bash +FORMS_DATABASE_URL=postgresql+psycopg://forms:forms@localhost:5432/forms +``` + +To try the service without running PostgreSQL: + +```bash +FORMS_DATABASE_URL=sqlite:///./forms.db +``` + +See [`.env.example`](.env.example) for every setting and its default. + +## Run + +Hymical Forms is two processes sharing one database. + +```bash +uvicorn hymical_forms.main:app --reload +``` + +```bash +python -m hymical_forms.worker +``` + +The **API** accepts submissions, stores them, and records that a webhook is owed. +It never makes an outbound request. The **worker** claims owed deliveries, sends +them, and retries the ones that fail. Running the API alone is fine: submissions +are still accepted and nothing is lost, they simply wait until a worker exists. + +Missing tables are created at startup, so an empty database is enough to begin. +Startup fails if the database cannot be reached, rather than serving requests +that would only fail later. There is no migration framework yet, so startup +never alters a table that already exists; see [Limitations](#limitations). + +> **Upgrading from an earlier build:** the schema has changed in every release so +> far, most recently by adding the `webhook_deliveries` table and giving +> `delivery_attempts` a `delivery_id` and `attempt_number`. Startup creates +> missing tables but never alters an existing one, so a database created before +> these changes has to be recreated. For local SQLite, delete the file and +> restart. For PostgreSQL, `DROP TABLE delivery_attempts, webhook_deliveries, +> submissions, endpoints;` and restart. There is no in-place upgrade path. +> +> Three consecutive schema changes with no migration tool is the clearest +> remaining infrastructure gap. Alembic is the next thing this project needs, +> and it should arrive before there is a database worth not dropping. + +Interactive API documentation is served at `http://127.0.0.1:8000/docs`. + +## API + +### `GET /health` + +Reports that the API process is running. + +```json +{ "status": "ok", "service": "hymical-forms", "version": "0.1.0" } +``` + +This is a liveness signal only. It does not check the database, so it stays +answerable while the database is down, which is what makes it useful for +deciding whether to restart the process. + +### `POST /endpoints` + +Registers a form endpoint. Submissions are only accepted for endpoints that +exist here. + +**This route is unauthenticated.** Authentication is deliberately out of scope +for now, so keep the service on a private network. + +```bash +curl -X POST http://127.0.0.1:8000/endpoints \ + -H 'Content-Type: application/json' \ + -d '{"id": "contact-form", "name": "Contact form", + "webhook_url": "https://example.com/hooks/forms"}' +``` + +| Field | Required | Meaning | +| ------------- | -------- | -------------------------------------------------- | +| `id` | yes | The public identifier the endpoint answers on | +| `name` | yes | Human-readable label, 1 to 200 characters | +| `is_active` | no | Whether it accepts submissions, defaults to `true` | +| `webhook_url` | no | Where accepted submissions are delivered | + +**Endpoint IDs** are supplied by you, not generated, because the ID appears in +the `action` URL of your HTML form and a memorable one is worth more than an +opaque one. An ID is 3 to 64 characters of lowercase ASCII letters, digits, `-` +and `_`, and must start and end with a letter or digit. It is also the primary +key, so it cannot be changed later. + +Returns `201 Created`: + +```json +{ + "id": "contact-form", + "name": "Contact form", + "is_active": true, + "created_at": "2026-08-24T14:34:27.432598Z", + "webhook_url": "https://example.com/hooks/forms", + "webhook_secret": "whsec_6f1c... (64 hex characters)" +} +``` + +> **Save `webhook_secret` now.** It is generated by the server, returned only in +> this response, and there is no route that reads it back. Losing it means +> creating a new endpoint. + +Reusing an ID returns `409 endpoint_already_exists`. There is no route to list, +update or delete endpoints yet, so a webhook can only be configured at creation +time. + +### `POST /f/{endpoint_id}` + +Accepts a form submission for a registered endpoint and stores it. + +A submission to an ID that does not exist is rejected with +`404 endpoint_not_found`, and one to an inactive endpoint with +`409 endpoint_inactive`. Neither leaves anything in the database. + +**Content types.** `application/x-www-form-urlencoded` and +`multipart/form-data` are both accepted, so a plain HTML ` +``` + +### Errors + +Every non-2xx response uses one envelope. `code` is stable and +machine-readable; `details` appears only when there is something concrete to +add. + +```json +{ + "error": { + "code": "too_many_fields", + "message": "Submission carries 120 fields, which exceeds the limit of 100.", + "details": { "limit": 100, "received": 120 } + } +} +``` + +| Status | `code` | Cause | +| ------ | -------------------------- | -------------------------------------------------------- | +| 400 | `malformed_form_body` | Body does not parse as the declared content type | +| 400 | `invalid_idempotency_key` | `Idempotency-Key` header breaks the key format rules | +| 404 | `invalid_endpoint_id` | Submission path is not a well-formed endpoint ID | +| 404 | `endpoint_not_found` | Endpoint ID is well formed but no such endpoint exists | +| 404 | `not_found` | Unknown path | +| 405 | `method_not_allowed` | Wrong method for a known path | +| 409 | `endpoint_inactive` | Endpoint exists but is not accepting submissions | +| 409 | `endpoint_already_exists` | Endpoint ID is already taken | +| 409 | `idempotency_conflict` | Idempotency key already used for different content | +| 413 | `request_body_too_large` | Body exceeded `FORMS_MAX_BODY_BYTES` | +| 415 | `unsupported_media_type` | Content type is not a supported form encoding | +| 422 | `empty_submission` | No fields were submitted | +| 422 | `invalid_endpoint_id` | Endpoint ID in a request body breaks the ID rules | +| 422 | `invalid_request` | Request body failed schema validation | +| 422 | `invalid_webhook_url` | Webhook destination is malformed or not permitted | +| 422 | `file_upload_not_supported`| A multipart part carried a file | +| 422 | ingestion rule codes | See below | +| 500 | `internal_error` | Unexpected failure; no internals are exposed | +| 503 | `storage_unavailable` | The database could not be reached or written to | + +Ingestion rule codes are `too_many_fields`, `field_name_too_long`, +`field_value_too_long`, `invalid_field_name` and `invalid_field_value`. + +`invalid_endpoint_id` carries a different status depending on where the ID came +from: `404` when it arrived as a submission path that addresses nothing, `422` +when it arrived as a field in a request body. + +## Configuration + +All settings are read from `FORMS_`-prefixed environment variables, or from a +`.env` file in the working directory. + +| Variable | Default | Meaning | +| ------------------------------ | ---------- | --------------------------------------- | +| `FORMS_DATABASE_URL` | *required* | SQLAlchemy database URL | +| `FORMS_MAX_BODY_BYTES` | `262144` | Largest accepted request body, in bytes | +| `FORMS_MAX_FIELDS` | `100` | Largest number of name/value pairs | +| `FORMS_MAX_FIELD_NAME_LENGTH` | `128` | Largest field name, in characters | +| `FORMS_MAX_FIELD_VALUE_LENGTH` | `16384` | Largest field value, in characters | +| `FORMS_WEBHOOK_CONNECT_TIMEOUT_SECONDS` | `5` | Wait for a webhook to accept a connection | +| `FORMS_WEBHOOK_READ_TIMEOUT_SECONDS` | `10` | Wait for a webhook to respond | +| `FORMS_ALLOW_PRIVATE_WEBHOOK_TARGETS` | `false` | Permit loopback and private webhook targets. Development only | +| `FORMS_WEBHOOK_MAX_ATTEMPTS` | `5` | Attempts before a delivery is given up on | +| `FORMS_WEBHOOK_RETRY_INITIAL_SECONDS` | `10` | Wait before the second attempt; later waits double | +| `FORMS_WEBHOOK_RETRY_MAX_SECONDS` | `3600` | Cap on the wait between attempts | +| `FORMS_WORKER_POLL_SECONDS` | `1` | How often an idle worker looks for work | +| `FORMS_WORKER_BATCH_SIZE` | `10` | Deliveries a worker claims at once | +| `FORMS_WORKER_LEASE_SECONDS` | `60` | How long a worker's claim holds | + +## Development + +```bash +pytest # run the test suite +ruff check . # lint +ruff format --check . # formatting check +mypy # type check +``` + +Tests run against an in-memory SQLite database, one per test, so no database +server is needed and nothing is left behind. + +### Layout + +``` +src/hymical_forms/ + app.py application assembly and startup + config.py typed settings + db.py engine, session, and schema lifecycle + errors.py the shared JSON error envelope + delivery.py the outbound webhook request itself + ingestion.py domain rules: endpoint IDs, submission validation + middleware.py request body size limit + models.py the persisted schema + storage.py queries and writes + webhooks.py webhook rules: URL validation, payload, signature, retry policy + worker.py the delivery worker process + main.py ASGI entrypoint + api/ HTTP routes and response models +``` + +`ingestion.py` and `webhooks.py` hold the domain rules and know nothing about +HTTP or the database. `models.py` and `storage.py` are the only modules that +write queries, and `delivery.py` is the only one that makes an outbound request. +`api/` translates requests into domain rules and storage calls, and their +outcomes into responses. `worker.py` is a separate process and shares only the +database with the API. + +### Storage notes + +Submission fields are stored as a JSON object mapping each field name to the +list of values submitted under it, which is how repeated names survive intact. +On PostgreSQL the column is `json` rather than `jsonb`, because `jsonb` +normalises object key order and would silently reorder a form's fields. + +Each request runs in one transaction, committed explicitly rather than in the +session teardown, so a failure becomes an error response instead of a success +for a row that never landed. A failure anywhere before the commit leaves the +database untouched. + +An idempotency key is unique per endpoint through a database constraint on +`(endpoint_id, idempotency_key)`. Both PostgreSQL and SQLite treat NULLs in a +unique constraint as distinct, so submissions sent without a key stay +unrestricted without needing a partial index. A lookup before inserting is only +an optimisation for the common retry; when two requests race, one insert loses +on the constraint, rolls back and reads the winner's row. A `CHECK` constraint +keeps the key and its fingerprint either both set or both absent. + +A submission and the delivery it owes are written in one transaction. Either +both land or neither does, so there is no state in which a form was accepted but +the promise to deliver it went missing, and none in which delivery work exists +for a submission that does not. + +The network call happens later, in the worker, with no transaction open. Holding +a database transaction across a call to somebody else's server would tie the +connection pool to how fast that server answers. + +Workers claim deliveries with `SELECT ... FOR UPDATE SKIP LOCKED` on PostgreSQL, +so two workers scanning at once are handed different rows rather than fighting +over the same one. SQLite has no such locking and silently ignores `FOR UPDATE`, +so the claim also performs a conditional update and treats a row as claimed only +if that update matched. That guard is redundant under `SKIP LOCKED` and is what +makes the claim safe on SQLite. + +## Limitations + +- **Delivery is at-least-once, never exactly-once.** A worker that delivers + successfully and dies before recording it will have its lease expire, and the + next worker will deliver the same event again. Deduplicate on the submission + `id` in the signed payload. +- **PostgreSQL worker concurrency is not exercised by the test suite.** Tests run + on SQLite, which cannot model `SELECT ... FOR UPDATE SKIP LOCKED`. The + generated PostgreSQL SQL is asserted, and the claim is written so that it is + also correct without row locking, but two real workers racing on PostgreSQL has + not been run. A PostgreSQL service in CI is the way to close this. +- **A failed delivery is final and cannot be replayed.** Once a delivery reaches + `failed`, nothing retries it and there is no manual replay route. +- **The lease must outlast a delivery attempt.** A batch is delivered + concurrently, so it takes about as long as its slowest single delivery rather + than the sum, but if `FORMS_WORKER_LEASE_SECONDS` were set below the connect + and read timeouts combined, another worker could claim a delivery that is still + in flight and send it twice. The defaults leave a wide margin; keep it that way + if you change them. +- **SSRF protection is partial.** Destination URLs are checked for scheme and for + literal internal addresses, and redirects are not followed. Hostnames are + **not** resolved, so a name that resolves to a private address still passes, + and DNS rebinding is not addressed at all. Closing this properly means + resolving at request time and pinning the connection to the validated address. + Treat the current checks as a guardrail against mistakes, not a defence against + an attacker who can configure endpoints. +- **No authentication.** Anyone who can reach the API can create an endpoint, + point its webhook anywhere, and post to any active one. There is no rate + limiting or spam protection. +- **A webhook can only be set when the endpoint is created.** There is no route + to change a destination or rotate a signing secret. +- **No API for delivery attempts.** They are recorded, but reading them means + querying the database directly. +- **No migration framework.** Startup creates missing tables and nothing else, + so any future change to an existing column has to be applied by hand. + Alembic will arrive when the schema first needs to change. +- **No way to read submissions back over the API.** They are stored, but + retrieval, export and retention are not implemented. +- **No route to list, update or delete endpoints.** +- **No file uploads.** Multipart text fields are accepted; file parts are + rejected. +- **`multipart/form-data` bodies are buffered in memory,** bounded by + `FORMS_MAX_BODY_BYTES`. +- A rejected submission reveals whether an endpoint ID exists, which allows + enumeration. This is unavoidable while the API is unauthenticated. +- **Idempotency keys never expire.** A key stays spent for as long as its + submission is stored, so the table only grows. Expiry belongs with retention. +- **Idempotency keys are shared across all clients of an endpoint,** because + there is nothing to scope them to yet. Guessing another client's key returns + that submission's ID and timestamp, though never its contents. Random keys of + the required length make this impractical, and API keys will close it properly. +- **A replay is only recognised once the first attempt has committed.** A retry + sent while the original is still in flight is treated as a concurrent request, + which is safe, but a retry sent after the original *failed* is a new + submission, which is correct. +- Submission IDs are opaque and not yet guaranteed stable in format. + +## License + +[Apache License 2.0](LICENSE). diff --git a/docs/images/logo_symbol_tramsparent.png b/docs/images/logo_symbol_transparent.png similarity index 100% rename from docs/images/logo_symbol_tramsparent.png rename to docs/images/logo_symbol_transparent.png diff --git a/pyproject.toml b/pyproject.toml new file mode 100644 index 0000000..84c4f77 --- /dev/null +++ b/pyproject.toml @@ -0,0 +1,77 @@ +[build-system] +requires = ["hatchling"] +build-backend = "hatchling.build" + +[project] +name = "hymical-forms" +dynamic = ["version"] +description = "Reliable form ingestion and webhook delivery for developers." +readme = "README.md" +license = "Apache-2.0" +license-files = ["LICENSE"] +requires-python = ">=3.11" +authors = [{ name = "Hymical" }] +keywords = ["forms", "webhooks", "ingestion", "fastapi"] +classifiers = [ + "Development Status :: 3 - Alpha", + "Framework :: FastAPI", + "Intended Audience :: Developers", + "Programming Language :: Python :: 3", + "Topic :: Internet :: WWW/HTTP :: HTTP Servers", +] +dependencies = [ + "fastapi>=0.115", + "httpx2>=2.0", # outbound webhook client, and the transport starlette's TestClient uses + "psycopg[binary]>=3.1", # PostgreSQL driver for the intended production database + "pydantic>=2.7", + "pydantic-settings>=2.3", + "python-multipart>=0.0.9", + "sqlalchemy>=2.0", + "uvicorn>=0.30", +] + +[project.optional-dependencies] +dev = [ + "mypy>=1.11", + "pytest>=8.3", + "ruff>=0.6", +] + +[project.urls] +Homepage = "https://github.com/hymical/forms" +Source = "https://github.com/hymical/forms" +Issues = "https://github.com/hymical/forms/issues" + +[tool.hatch.version] +path = "src/hymical_forms/__init__.py" + +[tool.hatch.build.targets.wheel] +packages = ["src/hymical_forms"] + +[tool.pytest.ini_options] +testpaths = ["tests"] +addopts = "-q --strict-markers --strict-config" + +[tool.ruff] +target-version = "py311" +line-length = 100 +src = ["src", "tests"] + +[tool.ruff.lint] +select = [ + "E", # pycodestyle errors + "F", # pyflakes + "I", # import sorting + "UP", # pyupgrade + "B", # bugbear + "SIM", # simplify + "RUF", # ruff-specific +] + +[tool.mypy] +python_version = "3.11" +files = ["src", "tests"] +strict = true +# Teaches mypy that a BaseSettings subclass can be constructed with no arguments +# because its values come from the environment. +plugins = ["pydantic.mypy"] diff --git a/src/hymical_forms/__init__.py b/src/hymical_forms/__init__.py new file mode 100644 index 0000000..1a58e53 --- /dev/null +++ b/src/hymical_forms/__init__.py @@ -0,0 +1,7 @@ +""" +hymical forms: reliable form ingestion and webhook delivery for developers +""" + +__version__ = "0.1.0" + +__all__ = ["__version__"] diff --git a/src/hymical_forms/api/__init__.py b/src/hymical_forms/api/__init__.py new file mode 100644 index 0000000..f0e63dd --- /dev/null +++ b/src/hymical_forms/api/__init__.py @@ -0,0 +1,3 @@ +""" +the HTTP layer: routing, request parsing, and response shapes +""" diff --git a/src/hymical_forms/api/endpoints.py b/src/hymical_forms/api/endpoints.py new file mode 100644 index 0000000..eb78fdd --- /dev/null +++ b/src/hymical_forms/api/endpoints.py @@ -0,0 +1,188 @@ +""" +endpoint management: ``POST /endpoints`` +""" + +from __future__ import annotations + +from datetime import datetime +from http import HTTPStatus + +from fastapi import APIRouter, Request +from pydantic import BaseModel, Field + +from hymical_forms import storage, webhooks +from hymical_forms.config import Settings +from hymical_forms.db import SessionDep +from hymical_forms.errors import ApiError, ErrorResponse +from hymical_forms.ingestion import ENDPOINT_ID_RULE, is_valid_endpoint_id +from hymical_forms.models import ENDPOINT_NAME_MAX_LENGTH +from hymical_forms.webhooks import WEBHOOK_URL_MAX_LENGTH + +router = APIRouter(tags=["endpoints"]) + + +class InvalidEndpointId(ApiError): + """ + raised when a submitted endpoint identifier does not follow the rules + """ + + # The ingestion route answers 404 for a malformed ID because the path simply + # does not address an endpoint. Here the ID arrives in a request body, where + # the same problem is an unprocessable field rather than a missing resource. + status_code = HTTPStatus.UNPROCESSABLE_ENTITY + code = "invalid_endpoint_id" + + def __init__(self) -> None: + """ + state the endpoint identifier rules the request failed + """ + super().__init__(ENDPOINT_ID_RULE, details={"field": "id"}) + + +class EndpointIdConflict(ApiError): + """ + raised when the requested endpoint identifier is already taken + """ + + status_code = HTTPStatus.CONFLICT + code = "endpoint_already_exists" + + def __init__(self, endpoint_id: str) -> None: + """ + name the endpoint identifier that is already in use + :param endpoint_id: the identifier that collided + """ + super().__init__( + f"An endpoint with the ID {endpoint_id!r} already exists.", + details={"endpoint_id": endpoint_id}, + ) + + +class InvalidWebhookUrl(ApiError): + """ + raised when a webhook destination is malformed or not permitted + """ + + status_code = HTTPStatus.UNPROCESSABLE_ENTITY + code = "invalid_webhook_url" + + def __init__(self, reason: str) -> None: + """ + report why the destination was refused + :param reason: short phrase completing "the webhook URL ..." + """ + super().__init__(f"The webhook URL {reason}.", details={"field": "webhook_url"}) + + +class CreateEndpointRequest(BaseModel): + """ + the body accepted when creating an endpoint + """ + + id: str = Field(description=ENDPOINT_ID_RULE) + name: str = Field( + min_length=1, + max_length=ENDPOINT_NAME_MAX_LENGTH, + description="Human-readable label, shown to whoever administers the endpoint.", + ) + is_active: bool = Field( + default=True, + description="Whether the endpoint accepts submissions. Inactive endpoints reject them.", + ) + webhook_url: str | None = Field( + default=None, + max_length=WEBHOOK_URL_MAX_LENGTH, + description=( + "Optional http or https destination to deliver accepted submissions to. " + "A signing secret is generated for it and returned once." + ), + ) + + +class EndpointResponse(BaseModel): + """ + an endpoint as returned by the API + """ + + id: str = Field(description="The public identifier the endpoint answers on.") + name: str = Field(description="Human-readable label for the endpoint.") + is_active: bool = Field(description="Whether the endpoint currently accepts submissions.") + created_at: datetime = Field(description="UTC timestamp of when the endpoint was created.") + webhook_url: str | None = Field( + description="Where accepted submissions are delivered, or null if none is configured." + ) + webhook_secret: str | None = Field( + description=( + "The signing secret for this endpoint's webhook. Returned only here, at " + "creation, and never retrievable again. Null if no webhook is configured." + ) + ) + + +@router.post( + "/endpoints", + status_code=HTTPStatus.CREATED, + summary="Create a form endpoint", + responses={ + 409: {"model": ErrorResponse, "description": "Endpoint ID already taken"}, + 422: { + "model": ErrorResponse, + "description": "Invalid endpoint ID, name, or webhook URL", + }, + 503: {"model": ErrorResponse, "description": "Database unavailable"}, + }, +) +def create_endpoint( + payload: CreateEndpointRequest, request: Request, session: SessionDep +) -> EndpointResponse: + """ + create an endpoint that submissions may then be addressed to + :param payload: the endpoint identifier, label, active state and optional webhook + :param request: the incoming request, read for the active configuration + :param session: the session this request does its database work through + :returns: the endpoint as persisted, including its signing secret if one was made + """ + # A plain ``def`` route, so FastAPI runs it in a worker thread and the + # synchronous database calls never block the event loop. + if not is_valid_endpoint_id(payload.id): + raise InvalidEndpointId() + + settings: Settings = request.app.state.settings + secret: str | None = None + if payload.webhook_url is not None: + try: + webhooks.validate_webhook_url( + payload.webhook_url, + allow_private_targets=settings.allow_private_webhook_targets, + ) + except webhooks.WebhookUrlRejected as exc: + raise InvalidWebhookUrl(exc.reason) from exc + # Generated here rather than accepted from the caller, so its strength is + # this service's responsibility and never an integrator's oversight. + secret = webhooks.new_signing_secret() + + try: + endpoint = storage.create_endpoint( + session, + endpoint_id=payload.id, + name=payload.name, + is_active=payload.is_active, + webhook_url=payload.webhook_url, + webhook_secret=secret, + ) + except storage.EndpointAlreadyExists as exc: + raise EndpointIdConflict(payload.id) from exc + + session.commit() + + # The secret leaves the service exactly once, in this response. There is no + # route that reads it back, so a caller that loses it has to make a new + # endpoint rather than being handed the old secret again. + return EndpointResponse( + id=endpoint.id, + name=endpoint.name, + is_active=endpoint.is_active, + created_at=endpoint.created_at, + webhook_url=endpoint.webhook_url, + webhook_secret=secret, + ) diff --git a/src/hymical_forms/api/health.py b/src/hymical_forms/api/health.py new file mode 100644 index 0000000..d3836ed --- /dev/null +++ b/src/hymical_forms/api/health.py @@ -0,0 +1,36 @@ +""" +health endpoint +""" + +from __future__ import annotations + +from typing import Literal + +from fastapi import APIRouter +from pydantic import BaseModel + +from hymical_forms import __version__ + +router = APIRouter(tags=["health"]) + + +class HealthResponse(BaseModel): + """ + liveness report for a hymical forms process + """ + + status: Literal["ok"] + service: str + version: str + + +@router.get("/health", summary="Report process health") +async def health() -> HealthResponse: + """ + report that the api process is running and able to serve requests + :returns: a liveness payload naming the service and its version + """ + # This is a liveness signal only. Hymical Forms has no external dependencies + # yet, so there is nothing to distinguish readiness from liveness; a separate + # readiness endpoint will arrive with persistence. + return HealthResponse(status="ok", service="hymical-forms", version=__version__) diff --git a/src/hymical_forms/api/submissions.py b/src/hymical_forms/api/submissions.py new file mode 100644 index 0000000..8276cb3 --- /dev/null +++ b/src/hymical_forms/api/submissions.py @@ -0,0 +1,433 @@ +""" +form ingestion endpoint: ``POST /f/{endpoint_id}`` +""" + +from __future__ import annotations + +import math +from datetime import datetime +from http import HTTPStatus + +from fastapi import APIRouter, Request +from pydantic import BaseModel, Field +from python_multipart.exceptions import ParseError +from starlette.concurrency import run_in_threadpool +from starlette.datastructures import UploadFile +from starlette.formparsers import FormParser, MultiPartException, MultiPartParser + +from hymical_forms import storage +from hymical_forms.config import Settings +from hymical_forms.db import SessionDep +from hymical_forms.errors import ApiError, ErrorResponse +from hymical_forms.ingestion import ( + ENDPOINT_ID_RULE, + IDEMPOTENCY_KEY_RULE, + build_submission, + is_valid_endpoint_id, + is_valid_idempotency_key, + payload_fingerprint, +) +from hymical_forms.webhooks import WebhookTarget + +IDEMPOTENCY_KEY_HEADER = "Idempotency-Key" + +URLENCODED = "application/x-www-form-urlencoded" +MULTIPART = "multipart/form-data" +SUPPORTED_MEDIA_TYPES = (URLENCODED, MULTIPART) + +# Content-Type values are echoed back to help developers debug their form tags, +# but only ever a bounded prefix of what the client sent. +_MEDIA_TYPE_ECHO_LIMIT = 128 + +router = APIRouter(tags=["submissions"]) + + +class InvalidEndpointId(ApiError): + """ + raised when the path segment is not a well-formed endpoint identifier + """ + + status_code = HTTPStatus.NOT_FOUND + code = "invalid_endpoint_id" + + def __init__(self) -> None: + """ + state the endpoint identifier rules the request failed + """ + super().__init__(f"The path does not address a form endpoint. {ENDPOINT_ID_RULE}") + + +class EndpointNotFound(ApiError): + """ + raised when the identifier is well formed but no such endpoint exists + """ + + # Deliberately the same status as a malformed ID: from outside, both mean the + # path does not address a form endpoint. The code tells the two apart. + status_code = HTTPStatus.NOT_FOUND + code = "endpoint_not_found" + + def __init__(self, endpoint_id: str) -> None: + """ + name the endpoint identifier that could not be resolved + :param endpoint_id: the identifier taken from the request path + """ + super().__init__( + f"No form endpoint with the ID {endpoint_id!r} exists.", + details={"endpoint_id": endpoint_id}, + ) + + +class EndpointInactive(ApiError): + """ + raised when the endpoint exists but is not accepting submissions + """ + + status_code = HTTPStatus.CONFLICT + code = "endpoint_inactive" + + def __init__(self, endpoint_id: str) -> None: + """ + name the endpoint identifier that is not accepting submissions + :param endpoint_id: the identifier taken from the request path + """ + super().__init__( + f"The form endpoint {endpoint_id!r} is not accepting submissions.", + details={"endpoint_id": endpoint_id}, + ) + + +class InvalidIdempotencyKey(ApiError): + """ + raised when the ``Idempotency-Key`` header is present but unusable + """ + + # A malformed header is a framing problem rather than a semantic one, which + # is what separates this from the 422 an unacceptable submission earns. + status_code = HTTPStatus.BAD_REQUEST + code = "invalid_idempotency_key" + + def __init__(self) -> None: + """ + state the idempotency key rules the request failed + """ + super().__init__( + f"The {IDEMPOTENCY_KEY_HEADER} header is not usable. {IDEMPOTENCY_KEY_RULE}" + ) + + +class IdempotencyConflict(ApiError): + """ + raised when an idempotency key was already spent on different content + """ + + status_code = HTTPStatus.CONFLICT + code = "idempotency_conflict" + + def __init__(self, endpoint_id: str, idempotency_key: str) -> None: + """ + report that the key is already tied to a different submission + :param endpoint_id: the endpoint the key is scoped to + :param idempotency_key: the key the client reused + """ + # The earlier submission's content is never described, only the fact that + # it differs, so the key cannot be used to read back someone else's form. + super().__init__( + f"The {IDEMPOTENCY_KEY_HEADER} {idempotency_key!r} was already used on endpoint " + f"{endpoint_id!r} for a different submission.", + details={"endpoint_id": endpoint_id, "idempotency_key": idempotency_key}, + ) + + +class UnsupportedMediaType(ApiError): + """ + raised when the request used a content type the endpoint cannot parse + """ + + status_code = HTTPStatus.UNSUPPORTED_MEDIA_TYPE + code = "unsupported_media_type" + + def __init__(self, received: str) -> None: + """ + report the rejected content type alongside the supported ones + :param received: the normalized media type taken from the request + """ + super().__init__( + f"Form submissions must be sent as {URLENCODED} or {MULTIPART}.", + details={ + "received": received[:_MEDIA_TYPE_ECHO_LIMIT] or None, + "supported": list(SUPPORTED_MEDIA_TYPES), + }, + ) + + +class MalformedFormBody(ApiError): + """ + raised when the body did not parse as the declared form content type + """ + + status_code = HTTPStatus.BAD_REQUEST + code = "malformed_form_body" + + def __init__(self, reason: str) -> None: + """ + report why the body could not be parsed + :param reason: the form parser's description of what went wrong + """ + super().__init__( + "The request body could not be parsed as form data.", + details={"reason": reason}, + ) + + +class FileUploadNotSupported(ApiError): + """ + raised when a multipart part carries a file, which this service does not accept + """ + + status_code = HTTPStatus.UNPROCESSABLE_ENTITY + code = "file_upload_not_supported" + + def __init__(self, field_name: str) -> None: + """ + name the field that carried a file part + :param field_name: name of the offending multipart field + """ + super().__init__( + f"Field {field_name!r} carries a file upload, which is not supported.", + details={"field": field_name}, + ) + + +class DeliveryStatus(BaseModel): + """ + whether this submission owes a webhook delivery + """ + + queued: bool = Field( + description=( + "True when a durable webhook delivery exists for this submission. False " + "when the endpoint has no webhook. A delivery is queued once and is not " + "queued again by an idempotent replay, so a replay of a webhook-enabled " + "submission still reports true." + ) + ) + + +class SubmissionAccepted(BaseModel): + """ + acknowledgement returned for an accepted submission + """ + + # The submitted values are not echoed back: the client already has them, and + # reflecting user input adds nothing but risk. + submission_id: str = Field(description="Opaque identifier generated for this submission.") + endpoint_id: str = Field(description="The endpoint the submission was addressed to.") + received_at: datetime = Field(description="UTC timestamp of when the API accepted the body.") + field_count: int = Field(description="Number of name/value pairs the submission carried.") + idempotent_replay: bool = Field( + description=( + "True when this response describes a submission an earlier request already " + "stored, rather than one created now. Always false without an " + "Idempotency-Key header." + ), + ) + delivery: DeliveryStatus = Field( + description="Whether a webhook delivery is owed for this submission." + ) + + +@router.post( + "/f/{endpoint_id}", + status_code=HTTPStatus.ACCEPTED, + summary="Submit a form", + responses={ + 400: {"model": ErrorResponse, "description": "Malformed form body or idempotency key"}, + 404: {"model": ErrorResponse, "description": "Invalid or unknown endpoint ID"}, + 409: { + "model": ErrorResponse, + "description": "Endpoint is not accepting submissions, or idempotency key reused", + }, + 413: {"model": ErrorResponse, "description": "Request body too large"}, + 415: {"model": ErrorResponse, "description": "Unsupported content type"}, + 422: {"model": ErrorResponse, "description": "Submission rejected by an ingestion rule"}, + 503: {"model": ErrorResponse, "description": "Database unavailable"}, + }, +) +async def submit(endpoint_id: str, request: Request, session: SessionDep) -> SubmissionAccepted: + """ + accept an html form submission and store it + :param endpoint_id: endpoint identifier taken from the request path + :param request: the incoming request, read for its content type and body + :param session: the session this request does its database work through + :returns: an acknowledgement carrying the stored submission's metadata + """ + # The response is 202 Accepted rather than 201 Created: the submission is + # stored, but the delivery it was accepted for has not happened yet. + # + # The endpoint is resolved before the body is parsed, so an unknown endpoint + # costs one indexed lookup rather than a full parse of a body we would throw + # away. This handler must stay ``async`` to stream the body, so each blocking + # database call is handed to a worker thread instead of stalling the loop. + if not is_valid_endpoint_id(endpoint_id): + raise InvalidEndpointId() + + endpoint = await run_in_threadpool(storage.get_endpoint, session, endpoint_id) + if endpoint is None: + raise EndpointNotFound(endpoint_id) + if not endpoint.is_active: + raise EndpointInactive(endpoint_id) + + # Read the webhook configuration off the row now, while the session is known + # to be clean. Storing the submission can roll back to settle an idempotency + # race, and a rollback expires loaded objects, so touching the endpoint later + # would silently issue a refresh query from this async handler. + webhook_url = endpoint.webhook_url + webhook_secret = endpoint.webhook_secret + + media_type = _media_type(request.headers.get("content-type")) + if media_type not in SUPPORTED_MEDIA_TYPES: + raise UnsupportedMediaType(media_type) + + idempotency_key = _idempotency_key(request) + + settings: Settings = request.app.state.settings + submission = build_submission( + endpoint_id, + await _parse_form(request, media_type, settings), + max_fields=settings.max_fields, + max_field_name_length=settings.max_field_name_length, + max_field_value_length=settings.max_field_value_length, + ) + + # The submission and, if the endpoint has a webhook, the durable obligation to + # deliver it are committed together. Nothing outbound happens here: once this + # returns, a worker owns the delivery, and a crash in this process can no + # longer lose a delivery that was implicitly promised by a 202. + # + # The commit happens inside the handler, not in the session dependency's + # teardown, so that a failure still becomes an error response. Teardown runs + # after the response has been sent, where raising could no longer change it. + webhook = ( + WebhookTarget(url=webhook_url, secret=webhook_secret) + if webhook_url is not None and webhook_secret is not None + else None + ) + try: + stored = await run_in_threadpool( + storage.store_submission, + session, + submission, + now=submission.received_at, + idempotency_key=idempotency_key, + payload_fingerprint=( + payload_fingerprint(submission.fields) if idempotency_key else None + ), + webhook=webhook, + ) + except storage.IdempotencyKeyReused as exc: + raise IdempotencyConflict(exc.endpoint_id, exc.idempotency_key) from exc + + # A replay answers with the original submission's identity and timestamp, so + # a client that retried after a lost response ends up describing one event. + # It reports the same queued state as the original, because the delivery that + # request created is still the one that is owed. + return SubmissionAccepted( + submission_id=stored.submission.id, + endpoint_id=stored.submission.endpoint_id, + received_at=stored.submission.received_at, + field_count=stored.submission.field_count, + idempotent_replay=stored.replayed, + delivery=DeliveryStatus(queued=webhook is not None), + ) + + +def _idempotency_key(request: Request) -> str | None: + """ + read and validate the retry key a client may have sent + :param request: the incoming request + :returns: the key, or None if the header was absent + :raises InvalidIdempotencyKey: if the header is present but breaks the key rules + """ + # An absent header keeps the pre-idempotency behaviour exactly: every accepted + # request stores a new submission. A header that is present but empty is a + # client bug, and treating it as absent would silently drop the guarantee the + # client was asking for. + key = request.headers.get(IDEMPOTENCY_KEY_HEADER) + if key is None: + return None + if not is_valid_idempotency_key(key): + raise InvalidIdempotencyKey() + return key + + +async def _parse_form( + request: Request, media_type: str, settings: Settings +) -> list[tuple[str, str]]: + """ + parse the body into ordered name/value pairs, preserving repeated names + :param request: the incoming request, streamed into the form parser + :param media_type: normalized media type taken from the Content-Type header + :param settings: active configuration, used to size the parser buffers + :returns: ordered name/value pairs exactly as submitted + :raises MalformedFormBody: if the body does not parse as the declared media type + :raises FileUploadNotSupported: if a multipart part carries a file + """ + # The parser is selected from the media type we normalized ourselves rather + # than through ``Request.form()``, whose dispatch compares the header verbatim + # even though media types are case-insensitive (RFC 9110 section 8.3). + # + # Starlette's own field and part limits are disabled: the request body size cap + # already bounds memory use, and leaving them on would let a library-defined + # threshold shadow the limits configured for this service. + parser: FormParser | MultiPartParser + if media_type == MULTIPART: + parser = MultiPartParser( + request.headers, + request.stream(), + max_files=math.inf, + max_fields=math.inf, + max_part_size=settings.max_body_bytes, + ) + else: + parser = FormParser( + request.headers, + request.stream(), + max_fields=math.inf, + max_part_size=settings.max_body_bytes, + ) + + try: + form = await parser.parse() + except (MultiPartException, ParseError) as exc: + raise MalformedFormBody(_failure_reason(exc)) from exc + + try: + items: list[tuple[str, str]] = [] + for name, value in form.multi_items(): + if isinstance(value, UploadFile): + raise FileUploadNotSupported(name) + items.append((name, value)) + return items + finally: + await form.close() + + +def _failure_reason(exc: MultiPartException | ParseError) -> str: + """ + extract a human-readable reason from a form parser failure + :param exc: the exception raised while parsing the body + :returns: the parser's description of what went wrong + """ + return exc.message if isinstance(exc, MultiPartException) else str(exc) + + +def _media_type(content_type: str | None) -> str: + """ + strip parameters such as charset and boundary from a Content-Type header + :param content_type: raw header value, or None when the header is absent + :returns: the lowercased media type, or an empty string when there is none + """ + if not content_type: + return "" + return content_type.split(";", 1)[0].strip().lower() diff --git a/src/hymical_forms/app.py b/src/hymical_forms/app.py new file mode 100644 index 0000000..ada7def --- /dev/null +++ b/src/hymical_forms/app.py @@ -0,0 +1,79 @@ +""" +application assembly +""" + +from __future__ import annotations + +from collections.abc import AsyncIterator +from contextlib import asynccontextmanager + +from fastapi import FastAPI + +from hymical_forms import __version__ +from hymical_forms.api import endpoints, health, submissions +from hymical_forms.config import Settings +from hymical_forms.db import create_engine_from_url, create_session_factory, init_db +from hymical_forms.errors import register_exception_handlers +from hymical_forms.middleware import BodySizeLimitMiddleware + +DESCRIPTION = """\ +Hymical Forms accepts HTML form submissions over HTTP so that developers do not +have to run a form backend of their own. + +Submissions are parsed, validated, and stored against a registered endpoint +together with the durable obligation to deliver them. A separate worker process +performs the webhook delivery and retries it. +""" + + +@asynccontextmanager +async def lifespan(app: FastAPI) -> AsyncIterator[None]: + """ + prepare the schema on startup and release the connection pool on shutdown + :param app: the application starting up + :returns: an async context manager wrapping the application's serving life + """ + # There is no migration framework yet, so creating missing tables at startup + # is the whole schema story. It is safe to repeat and never alters a table + # that already exists, which also means a changed column needs manual work. + init_db(app.state.engine) + yield + app.state.engine.dispose() + + +def create_app(settings: Settings | None = None) -> FastAPI: + """ + build a hymical forms application + :param settings: configuration to use, or None to load it from the environment + :returns: the configured FastAPI application + """ + # Settings, the engine and the session factory are attached to ``app.state`` + # rather than read from module-level singletons, so a test (or a future + # multi-tenant host) can run several differently configured applications in + # one process, each with its own database. + settings = settings or Settings() + + app = FastAPI( + title="Hymical Forms", + summary="Reliable form ingestion and webhook delivery for developers.", + description=DESCRIPTION, + version=__version__, + license_info={"name": "Apache-2.0", "identifier": "Apache-2.0"}, + lifespan=lifespan, + ) + engine = create_engine_from_url(settings.database_url) + app.state.settings = settings + app.state.engine = engine + app.state.session_factory = create_session_factory(engine) + # The API process holds no outbound HTTP client. Webhook delivery belongs to + # the worker, and having nothing to send with is the plainest way to keep the + # ingestion path free of network calls. + + app.add_middleware(BodySizeLimitMiddleware, max_bytes=settings.max_body_bytes) + register_exception_handlers(app) + + app.include_router(health.router) + app.include_router(endpoints.router) + app.include_router(submissions.router) + + return app diff --git a/src/hymical_forms/config.py b/src/hymical_forms/config.py new file mode 100644 index 0000000..71482b6 --- /dev/null +++ b/src/hymical_forms/config.py @@ -0,0 +1,117 @@ +""" +application settings, read from ``FORMS_``-prefixed environment variables + +Settings are added only when the code actually uses them, so this model is +currently limited to the ingestion boundary's protective limits. +""" + +from __future__ import annotations + +from pydantic import Field +from pydantic_settings import BaseSettings, SettingsConfigDict + +from hymical_forms.webhooks import RetryPolicy + + +class Settings(BaseSettings): + """ + runtime configuration for a hymical forms process + """ + + model_config = SettingsConfigDict( + env_prefix="FORMS_", + env_file=".env", + env_file_encoding="utf-8", + extra="ignore", + frozen=True, + ) + + database_url: str = Field( + description=( + "SQLAlchemy database URL. PostgreSQL is the intended production database, " + "for example postgresql+psycopg://user:password@localhost:5432/forms." + ), + ) + max_body_bytes: int = Field( + default=256 * 1024, + ge=1, + description="Largest request body accepted, in bytes. File uploads are not supported.", + ) + max_fields: int = Field( + default=100, + ge=1, + description="Largest number of name/value pairs accepted in one submission.", + ) + max_field_name_length: int = Field( + default=128, + ge=1, + description="Largest field name accepted, in characters.", + ) + max_field_value_length: int = Field( + default=16 * 1024, + ge=1, + description="Largest field value accepted, in characters.", + ) + webhook_connect_timeout_seconds: float = Field( + default=5.0, + gt=0, + description="How long to wait for a webhook destination to accept a connection.", + ) + webhook_read_timeout_seconds: float = Field( + default=10.0, + gt=0, + description="How long to wait for a webhook destination to respond.", + ) + webhook_max_attempts: int = Field( + default=5, + ge=1, + description="How many delivery attempts a submission gets before it is given up on.", + ) + webhook_retry_initial_seconds: float = Field( + default=10.0, + gt=0, + description="Wait before the second attempt. Each later wait doubles it.", + ) + webhook_retry_max_seconds: float = Field( + default=3600.0, + gt=0, + description="Cap on the wait between attempts, however far the backoff has doubled.", + ) + worker_poll_seconds: float = Field( + default=1.0, + gt=0, + description="How long the worker waits before looking for due deliveries again.", + ) + worker_batch_size: int = Field( + default=10, + ge=1, + description="How many deliveries a worker claims at once.", + ) + worker_lease_seconds: float = Field( + default=60.0, + gt=0, + description=( + "How long a worker's claim on a delivery holds. After this the delivery " + "becomes claimable again, which is how work is recovered from a worker " + "that died holding it." + ), + ) + allow_private_webhook_targets: bool = Field( + default=False, + description=( + "Permit webhook destinations on loopback and private addresses. " + "For local development and tests only; enabling it in production " + "exposes the server to SSRF." + ), + ) + + def retry_policy(self) -> RetryPolicy: + """ + gather the retry settings into the value the delivery code works with + :returns: the configured retry policy + """ + return RetryPolicy( + max_attempts=self.webhook_max_attempts, + initial_seconds=self.webhook_retry_initial_seconds, + max_seconds=self.webhook_retry_max_seconds, + ) diff --git a/src/hymical_forms/db.py b/src/hymical_forms/db.py new file mode 100644 index 0000000..3719387 --- /dev/null +++ b/src/hymical_forms/db.py @@ -0,0 +1,90 @@ +""" +database engine, session, and schema lifecycle +""" + +from __future__ import annotations + +from collections.abc import Iterator +from typing import Annotated, Any + +from fastapi import Depends +from sqlalchemy import Engine, create_engine, event +from sqlalchemy.engine import make_url +from sqlalchemy.orm import Session, sessionmaker +from sqlalchemy.pool import StaticPool +from starlette.requests import Request + +from hymical_forms.models import Base + + +def create_engine_from_url(url: str) -> Engine: + """ + build the engine for a database URL + :param url: SQLAlchemy database URL, such as ``postgresql+psycopg://.../forms`` + :returns: an engine configured for that backend + """ + parsed = make_url(url) + if parsed.get_backend_name() != "sqlite": + return create_engine(url) + + # SQLite backs the test suite and local experimentation, never production. + # Requests are served from a thread pool so a connection is not pinned to the + # thread that opened it, and an in-memory database only exists for as long as + # its single connection is held, which is what StaticPool guarantees. + kwargs: dict[str, Any] = {"connect_args": {"check_same_thread": False}} + if parsed.database in (None, "", ":memory:"): + kwargs["poolclass"] = StaticPool + + engine = create_engine(url, **kwargs) + event.listen(engine, "connect", _enable_sqlite_foreign_keys) + return engine + + +def _enable_sqlite_foreign_keys(dbapi_connection: Any, connection_record: Any) -> None: + """ + switch on SQLite foreign key enforcement, which is off by default + :param dbapi_connection: the freshly opened DBAPI connection + :param connection_record: the pool's bookkeeping record, unused + """ + # Without this, SQLite accepts rows PostgreSQL would reject, and the test + # suite would stop being a faithful stand-in for the real database. + cursor = dbapi_connection.cursor() + cursor.execute("PRAGMA foreign_keys=ON") + cursor.close() + + +def create_session_factory(engine: Engine) -> sessionmaker[Session]: + """ + build the session factory an application will serve requests from + :param engine: the engine sessions should be bound to + :returns: a configured session factory + """ + # ``expire_on_commit=False`` keeps loaded values readable after a commit, so + # building a response out of a just-committed row costs no extra query. + return sessionmaker(bind=engine, expire_on_commit=False) + + +def init_db(engine: Engine) -> None: + """ + create any tables that do not exist yet + :param engine: the engine whose database should hold the schema + """ + # There is no migration framework yet, so this is the whole schema story: it + # creates missing tables and never alters existing ones. + Base.metadata.create_all(engine) + + +def get_session(request: Request) -> Iterator[Session]: + """ + provide the session a request should do its database work through + :param request: the request being served + :returns: an iterator yielding one session, closed when the request ends + """ + # Closing a session rolls back whatever was not committed, so a handler that + # raises part way through leaves nothing behind. + factory: sessionmaker[Session] = request.app.state.session_factory + with factory() as session: + yield session + + +SessionDep = Annotated[Session, Depends(get_session)] diff --git a/src/hymical_forms/delivery.py b/src/hymical_forms/delivery.py new file mode 100644 index 0000000..fec5ad0 --- /dev/null +++ b/src/hymical_forms/delivery.py @@ -0,0 +1,107 @@ +""" +the outbound half of webhook delivery: one HTTP attempt, no retries + +This module is the only place the service makes an outbound request. It is kept +separate from :mod:`hymical_forms.webhooks` so that the rules and the payload can +be tested without a network, and separate from the database layer so that no +transaction is ever held open across a call to somebody else's server. +""" + +from __future__ import annotations + +import httpx2 + +from hymical_forms import __version__ +from hymical_forms.config import Settings +from hymical_forms.webhooks import ( + DELIVERY_ERROR_MAX_LENGTH, + SIGNATURE_HEADER, + DeliveryOutcome, + DeliveryResult, + sign, +) + +USER_AGENT = f"Hymical-Forms/{__version__}" + + +def create_webhook_client(settings: Settings) -> httpx2.AsyncClient: + """ + build the client every outbound webhook is sent through + :param settings: active configuration, read for its timeouts + :returns: a client with explicit timeouts, no redirects and no retries + """ + # One client for the process, so connections are reused rather than + # renegotiated per submission. + # + # Redirects are not followed. Beyond being surprising for a webhook, following + # them would let a destination bounce the request to an address that the URL + # validation refused, which is the usual way SSRF protection gets walked + # around. A 3xx is therefore reported as an unsuccessful HTTP status. + # + # Transport retries are pinned to zero. This interval promises exactly one + # attempt per submission, and a client that quietly retried would break that + # promise for any receiver that is not idempotent. + return httpx2.AsyncClient( + timeout=httpx2.Timeout( + connect=settings.webhook_connect_timeout_seconds, + read=settings.webhook_read_timeout_seconds, + write=settings.webhook_read_timeout_seconds, + pool=settings.webhook_connect_timeout_seconds, + ), + follow_redirects=False, + transport=httpx2.AsyncHTTPTransport(retries=0), + ) + + +async def deliver( + client: httpx2.AsyncClient, *, url: str, secret: str, body: bytes +) -> DeliveryResult: + """ + make one attempt to deliver a signed payload + :param client: the shared outbound client + :param url: the destination to post to + :param secret: the destination's signing secret + :param body: the exact bytes to sign and transmit + :returns: the outcome, never raising for a destination that misbehaves + """ + # ``content=body`` transmits these bytes verbatim. Passing the payload object + # and letting the client serialize it would sign one encoding and send + # another, and the receiver's signature check would fail for reasons nobody + # could see. + headers = { + "Content-Type": "application/json", + SIGNATURE_HEADER: sign(body, secret), + "User-Agent": USER_AGENT, + } + + try: + response = await client.post(url, content=body, headers=headers) + except httpx2.TimeoutException as exc: + return DeliveryResult(DeliveryOutcome.TIMEOUT, error=_describe("timed out", exc)) + except httpx2.RequestError as exc: + return DeliveryResult( + DeliveryOutcome.NETWORK_ERROR, error=_describe("could not connect", exc) + ) + + if 200 <= response.status_code < 300: + return DeliveryResult(DeliveryOutcome.SUCCEEDED, response_status=response.status_code) + + return DeliveryResult( + DeliveryOutcome.HTTP_ERROR, + response_status=response.status_code, + error=f"destination responded with HTTP {response.status_code}", + ) + + +def _describe(summary: str, exc: Exception) -> str: + """ + build bounded failure text for a delivery that never reached a response + :param summary: our own words for what went wrong + :param exc: the transport error raised while trying + :returns: a short message safe to store + """ + # The exception's own text is useful to whoever debugs this later, but it is + # shaped by the destination, so it is truncated before it reaches a column. + detail = str(exc).strip() + message = f"{summary}: {detail}" if detail else summary + return message[:DELIVERY_ERROR_MAX_LENGTH] diff --git a/src/hymical_forms/errors.py b/src/hymical_forms/errors.py new file mode 100644 index 0000000..e077747 --- /dev/null +++ b/src/hymical_forms/errors.py @@ -0,0 +1,232 @@ +""" +the single JSON error envelope used by every non-2xx response + +Every error the API can produce, whether raised by our own code, by FastAPI's +request validation, or by Starlette's routing, is rendered as:: + + {"error": {"code": "...", "message": "...", "details": {...}}} + +``code`` is a stable, machine-readable string; ``message`` is a human-readable +sentence; ``details`` is present only when there is something concrete to add +(the limit that was exceeded, the field at fault). Nothing in the envelope +exposes internal types, stack frames, or file paths. +""" + +from __future__ import annotations + +from http import HTTPStatus +from typing import Any, ClassVar + +from fastapi import FastAPI +from fastapi.exceptions import RequestValidationError +from pydantic import BaseModel, Field +from sqlalchemy.exc import SQLAlchemyError +from starlette.exceptions import HTTPException as StarletteHTTPException +from starlette.requests import Request +from starlette.responses import JSONResponse, Response + +from hymical_forms.ingestion import SubmissionRejected + + +class ErrorDetail(BaseModel): + """ + the body of an error response + """ + + code: str = Field(description="Stable, machine-readable error identifier.") + message: str = Field(description="Human-readable explanation of the failure.") + details: dict[str, Any] | None = Field( + default=None, + description="Optional structured context, such as the limit that was exceeded.", + ) + + +class ErrorResponse(BaseModel): + """ + the envelope returned for every error + """ + + error: ErrorDetail + + +class ApiError(Exception): + """ + an error that maps directly onto the public error envelope + """ + + # Subclasses fix ``status_code`` and ``code``; instances supply the message + # and any structured details. + status_code: ClassVar[int] = 500 + code: ClassVar[str] = "internal_error" + + def __init__(self, message: str, *, details: dict[str, Any] | None = None) -> None: + """ + record the message and context for an error response + :param message: human-readable explanation of the failure + :param details: optional structured context, such as the limit that was exceeded + """ + super().__init__(message) + self.message = message + self.details = details + + def as_response(self) -> JSONResponse: + """ + render this error in the shared envelope + :returns: a JSONResponse carrying the envelope and this error's status code + """ + return error_response( + status_code=self.status_code, + code=self.code, + message=self.message, + details=self.details, + ) + + +def error_response( + *, + status_code: int, + code: str, + message: str, + details: dict[str, Any] | None = None, +) -> JSONResponse: + """ + build a JSON response in the standard error envelope + :param status_code: HTTP status code to return + :param code: stable, machine-readable error identifier + :param message: human-readable explanation of the failure + :param details: optional structured context, omitted from the body when absent + :returns: a JSONResponse carrying the envelope + """ + payload = ErrorResponse(error=ErrorDetail(code=code, message=message, details=details)) + return JSONResponse(status_code=status_code, content=payload.model_dump(exclude_none=True)) + + +def register_exception_handlers(app: FastAPI) -> None: + """ + route every error class the app can raise through the shared envelope + :param app: the application to register the handlers on + """ + app.add_exception_handler(ApiError, _handle_api_error) + app.add_exception_handler(SubmissionRejected, _handle_submission_rejected) + app.add_exception_handler(SQLAlchemyError, _handle_storage_error) + app.add_exception_handler(StarletteHTTPException, _handle_http_exception) + app.add_exception_handler(RequestValidationError, _handle_request_validation_error) + app.add_exception_handler(Exception, _handle_unexpected_error) + + +# Starlette types every handler as ``(Request, Exception) -> Response``, so each +# handler re-narrows the exception it was registered for. + + +async def _handle_api_error(request: Request, exc: Exception) -> Response: + """ + render an error raised by our own HTTP layer + :param request: the request being handled + :param exc: the raised exception, always an ApiError + :returns: the envelope response + """ + assert isinstance(exc, ApiError) + return exc.as_response() + + +async def _handle_submission_rejected(request: Request, exc: Exception) -> Response: + """ + render a domain rejection + :param request: the request being handled + :param exc: the raised exception, always a SubmissionRejected + :returns: the envelope response, with a 422 status + """ + # Every ingestion rule failure is a well-formed request carrying an + # unacceptable submission, which is exactly what 422 describes. + assert isinstance(exc, SubmissionRejected) + return error_response( + status_code=HTTPStatus.UNPROCESSABLE_ENTITY, + code=exc.code, + message=exc.message, + details=exc.details, + ) + + +async def _handle_storage_error(request: Request, exc: Exception) -> Response: + """ + render a database failure without describing it + :param request: the request being handled + :param exc: the raised exception, always a SQLAlchemyError + :returns: the envelope response, with a 503 status + """ + # Driver messages carry table names, SQL text and sometimes connection + # details, so none of the exception reaches the client. 503 rather than 500 + # because the request itself was fine and retrying it may well succeed. + return error_response( + status_code=HTTPStatus.SERVICE_UNAVAILABLE, + code="storage_unavailable", + message="The submission could not be stored. Try again shortly.", + ) + + +async def _handle_http_exception(request: Request, exc: Exception) -> Response: + """ + render routing-level errors such as unknown paths and wrong methods + :param request: the request being handled + :param exc: the raised exception, always a Starlette HTTPException + :returns: the envelope response + """ + assert isinstance(exc, StarletteHTTPException) + return error_response( + status_code=exc.status_code, + code=_code_for_status(exc.status_code), + message=str(exc.detail), + ) + + +async def _handle_request_validation_error(request: Request, exc: Exception) -> Response: + """ + render a request that FastAPI could not validate + :param request: the request being handled + :param exc: the raised exception, always a RequestValidationError + :returns: the envelope response, with a 422 status + """ + assert isinstance(exc, RequestValidationError) + # Only the location and pydantic's short explanation are relayed. The raw + # error carries the offending input, which may be user data we should not + # reflect back, and internal type names that mean nothing to a caller. + fields = [ + { + "field": ".".join(str(part) for part in error["loc"][1:]) or None, + "issue": error["msg"], + } + for error in exc.errors() + ] + return error_response( + status_code=HTTPStatus.UNPROCESSABLE_ENTITY, + code="invalid_request", + message="The request body could not be validated.", + details={"fields": fields} if fields else None, + ) + + +async def _handle_unexpected_error(request: Request, exc: Exception) -> Response: + """ + return an opaque 500 rather than letting an internal error reach the client + :param request: the request being handled + :param exc: the unhandled exception, deliberately not described to the client + :returns: the envelope response, with a 500 status + """ + return error_response( + status_code=HTTPStatus.INTERNAL_SERVER_ERROR, + code="internal_error", + message="The request could not be processed.", + ) + + +def _code_for_status(status_code: int) -> str: + """ + derive an error code from a status code, so that 405 gives ``method_not_allowed`` + :param status_code: HTTP status code to name + :returns: the status phrase in snake case, or ``http_error`` if unrecognised + """ + try: + phrase = HTTPStatus(status_code).phrase + except ValueError: + return "http_error" + return phrase.lower().replace("-", " ").replace(" ", "_") diff --git a/src/hymical_forms/ingestion.py b/src/hymical_forms/ingestion.py new file mode 100644 index 0000000..872add1 --- /dev/null +++ b/src/hymical_forms/ingestion.py @@ -0,0 +1,251 @@ +""" +ingestion domain rules: endpoint identifiers and submission normalization + +This module is deliberately free of HTTP concepts. It answers two questions, +"is this a well-formed endpoint identifier?" and "is this set of name/value +pairs an acceptable submission?", and leaves status codes and wire formats to +the API layer. +""" + +from __future__ import annotations + +import hashlib +import json +import re +import uuid +from collections.abc import Mapping, Sequence +from dataclasses import dataclass +from datetime import UTC, datetime +from typing import Any + +ENDPOINT_ID_MIN_LENGTH = 3 +ENDPOINT_ID_MAX_LENGTH = 64 + +# Lowercase only, so that an endpoint ID has exactly one spelling. Hyphen and +# underscore are allowed inside, but an ID may not start or end with them. +_ENDPOINT_ID_PATTERN = re.compile(r"[a-z0-9](?:[a-z0-9_-]*[a-z0-9])?") + +# C0/C1 controls and DEL. HTML permits almost anything else in a field name +# (``user[email]``, ``entry.42``, non-ASCII labels), so nothing else is rejected. +_CONTROL_CHARS = re.compile(r"[\x00-\x1f\x7f-\x9f]") + +SUBMISSION_ID_PREFIX = "sub_" + +# ``sub_`` followed by a uuid4 in hex. +SUBMISSION_ID_MAX_LENGTH = len(SUBMISSION_ID_PREFIX) + 32 + +# Stated once so that every error mentioning the rule words it identically. +ENDPOINT_ID_RULE = ( + f"Endpoint IDs are {ENDPOINT_ID_MIN_LENGTH}-{ENDPOINT_ID_MAX_LENGTH} characters using " + "lowercase letters, digits, '-' and '_', and must start and end with a letter or digit." +) + +# An idempotency key is scoped to one endpoint, and this API is unauthenticated, +# so every client of an endpoint draws from the same key space. A short or +# predictable key would therefore collide with a stranger's submission, which is +# why the floor is high enough to force a random token rather than a counter. +IDEMPOTENCY_KEY_MIN_LENGTH = 16 +IDEMPOTENCY_KEY_MAX_LENGTH = 255 + +# Printable ASCII with no spaces: covers UUIDs, hex, base64 and base64url, and +# keeps unbounded or unprintable header content out of the database. +_IDEMPOTENCY_KEY_PATTERN = re.compile(r"[!-~]+") + +IDEMPOTENCY_KEY_RULE = ( + f"Idempotency keys are {IDEMPOTENCY_KEY_MIN_LENGTH}-{IDEMPOTENCY_KEY_MAX_LENGTH} printable " + "ASCII characters with no spaces. Use a random value such as a UUID." +) + +# A SHA-256 digest rendered as hex. +PAYLOAD_FINGERPRINT_LENGTH = 64 + + +def is_valid_idempotency_key(value: str) -> bool: + """ + report whether a header value is a usable idempotency key + :param value: the raw ``Idempotency-Key`` header value + :returns: True if the key is well formed + """ + return ( + IDEMPOTENCY_KEY_MIN_LENGTH <= len(value) <= IDEMPOTENCY_KEY_MAX_LENGTH + and _IDEMPOTENCY_KEY_PATTERN.fullmatch(value) is not None + ) + + +def payload_fingerprint(fields: Mapping[str, tuple[str, ...]]) -> str: + """ + digest the submitted content so a retry can be recognised as the same request + :param fields: the normalized fields of a submission + :returns: a hex SHA-256 digest of the field content + """ + # Only the fields are hashed. The generated submission ID and the received + # timestamp differ on every attempt, so taking the mapping rather than the + # whole submission makes their exclusion structural instead of a promise. + # + # The canonical form is a JSON array of ``[name, [values]]`` pairs, so it is + # sensitive to field order and to repeated values, both of which this service + # already treats as meaningful. JSON also makes the framing unambiguous: + # concatenating names and values would let two different submissions produce + # identical bytes. Nothing here depends on Python's randomized hashing, so + # the digest is stable across processes and restarts. + canonical = json.dumps( + [[name, list(values)] for name, values in fields.items()], + separators=(",", ":"), + ensure_ascii=False, + ) + return hashlib.sha256(canonical.encode("utf-8")).hexdigest() + + +def is_valid_endpoint_id(value: str) -> bool: + """ + report whether a path segment is a syntactically valid endpoint identifier + :param value: candidate identifier taken from the request path + :returns: True if the identifier is well formed + """ + # Interval 1 has no endpoint registry, so any syntactically valid identifier + # is treated as addressable. + return ( + ENDPOINT_ID_MIN_LENGTH <= len(value) <= ENDPOINT_ID_MAX_LENGTH + and _ENDPOINT_ID_PATTERN.fullmatch(value) is not None + ) + + +class SubmissionRejected(Exception): + """ + raised when a submission violates an ingestion rule and must not be accepted + """ + + def __init__(self, code: str, message: str, details: dict[str, Any] | None = None) -> None: + """ + record why a submission was refused + :param code: stable, machine-readable identifier for the broken rule + :param message: human-readable explanation of the failure + :param details: optional structured context, such as the limit that was exceeded + """ + super().__init__(message) + self.code = code + self.message = message + self.details = details + + +@dataclass(frozen=True, slots=True) +class Submission: + """ + a validated form submission in its internal representation + """ + + # Repeated field names are preserved as ordered tuples because HTML forms use + # them for checkbox groups and multi-selects; collapsing them would silently + # discard user input. + id: str + endpoint_id: str + received_at: datetime + fields: dict[str, tuple[str, ...]] + + @property + def field_count(self) -> int: + """ + count the name/value pairs the submission carries + :returns: the total number of submitted values across all field names + """ + return sum(len(values) for values in self.fields.values()) + + +def new_submission_id() -> str: + """ + generate an opaque, prefixed submission identifier + :returns: a fresh submission id such as ``sub_1f0c9a...`` + """ + return f"{SUBMISSION_ID_PREFIX}{uuid.uuid4().hex}" + + +def build_submission( + endpoint_id: str, + items: Sequence[tuple[str, str]], + *, + max_fields: int, + max_field_name_length: int, + max_field_value_length: int, +) -> Submission: + """ + validate parsed form pairs and normalize them into a submission + :param endpoint_id: the endpoint the submission was addressed to + :param items: ordered name/value pairs as parsed from the request body, repeats included + :param max_fields: largest number of name/value pairs accepted + :param max_field_name_length: largest field name accepted, in characters + :param max_field_value_length: largest field value accepted, in characters + :returns: the normalized submission + :raises SubmissionRejected: if the submission is empty or breaches a limit + """ + if len(items) > max_fields: + raise SubmissionRejected( + "too_many_fields", + f"Submission carries {len(items)} fields, which exceeds the limit of {max_fields}.", + {"limit": max_fields, "received": len(items)}, + ) + + fields: dict[str, tuple[str, ...]] = {} + for name, value in items: + _validate_field_name(name, max_field_name_length) + _validate_field_value(name, value, max_field_value_length) + fields[name] = (*fields.get(name, ()), value) + + if not fields: + raise SubmissionRejected( + "empty_submission", + "Submission contains no fields.", + ) + + return Submission( + id=new_submission_id(), + endpoint_id=endpoint_id, + received_at=datetime.now(UTC), + fields=fields, + ) + + +def _validate_field_name(name: str, max_length: int) -> None: + """ + check a submitted field name against the name rules + :param name: field name as submitted + :param max_length: largest field name accepted, in characters + :raises SubmissionRejected: if the name is empty, too long, or holds control characters + """ + if not name: + raise SubmissionRejected( + "invalid_field_name", + "Submission contains a field with an empty name.", + ) + if len(name) > max_length: + raise SubmissionRejected( + "field_name_too_long", + f"A field name exceeds the limit of {max_length} characters.", + {"limit": max_length, "received": len(name)}, + ) + if _CONTROL_CHARS.search(name): + raise SubmissionRejected( + "invalid_field_name", + "Submission contains a field name with control characters.", + ) + + +def _validate_field_value(name: str, value: str, max_length: int) -> None: + """ + check a submitted field value against the value rules + :param name: field name the value belongs to, used only in the error message + :param value: field value as submitted + :param max_length: largest field value accepted, in characters + :raises SubmissionRejected: if the value is too long or holds a null byte + """ + if len(value) > max_length: + raise SubmissionRejected( + "field_value_too_long", + f"The value of field {name!r} exceeds the limit of {max_length} characters.", + {"field": name, "limit": max_length, "received": len(value)}, + ) + if "\x00" in value: + raise SubmissionRejected( + "invalid_field_value", + f"The value of field {name!r} contains a null byte.", + {"field": name}, + ) diff --git a/src/hymical_forms/main.py b/src/hymical_forms/main.py new file mode 100644 index 0000000..98904c4 --- /dev/null +++ b/src/hymical_forms/main.py @@ -0,0 +1,13 @@ +""" +the ASGI entrypoint + +Run with:: + + uvicorn hymical_forms.main:app +""" + +from __future__ import annotations + +from hymical_forms.app import create_app + +app = create_app() diff --git a/src/hymical_forms/middleware.py b/src/hymical_forms/middleware.py new file mode 100644 index 0000000..01fbeb7 --- /dev/null +++ b/src/hymical_forms/middleware.py @@ -0,0 +1,101 @@ +""" +middleware protecting the ingestion boundary at the ASGI layer +""" + +from __future__ import annotations + +from http import HTTPStatus + +from starlette.types import ASGIApp, Message, Receive, Scope, Send + +from hymical_forms.errors import ApiError + + +class RequestBodyTooLarge(ApiError): + """ + raised when a request body exceeds the configured maximum + """ + + status_code = HTTPStatus.REQUEST_ENTITY_TOO_LARGE + code = "request_body_too_large" + + def __init__(self, limit: int) -> None: + """ + record the limit the body overran + :param limit: largest request body accepted, in bytes + """ + super().__init__( + f"Request body exceeds the limit of {limit} bytes.", + details={"limit_bytes": limit}, + ) + + +class BodySizeLimitMiddleware: + """ + reject requests whose body exceeds a configured size + """ + + # Starlette buffers request bodies without an upper bound, so the cap has to + # sit in front of the form parsers rather than inside a route handler. + + def __init__(self, app: ASGIApp, *, max_bytes: int) -> None: + """ + wrap an ASGI application with a request body size cap + :param app: the ASGI application to wrap + :param max_bytes: largest request body accepted, in bytes + """ + self.app = app + self.max_bytes = max_bytes + + async def __call__(self, scope: Scope, receive: Receive, send: Send) -> None: + """ + pass the request through, refusing any body over the limit + :param scope: ASGI connection scope + :param receive: ASGI callable yielding request messages + :param send: ASGI callable accepting response messages + """ + if scope["type"] != "http": + await self.app(scope, receive, send) + return + + # A request that declares an oversized Content-Length is refused before a + # single body byte is read. + declared = _declared_content_length(scope) + if declared is not None and declared > self.max_bytes: + await RequestBodyTooLarge(self.max_bytes).as_response()(scope, receive, send) + return + + received = 0 + + async def limited_receive() -> Message: + """ + read the next request message, cutting off an oversized body + :returns: the next ASGI message + :raises RequestBodyTooLarge: once the running body total crosses the limit + """ + nonlocal received + message = await receive() + if message["type"] == "http.request": + received += len(message.get("body", b"")) + if received > self.max_bytes: + # Raised inside the application, so the registered ApiError + # handler renders it in the standard envelope. + raise RequestBodyTooLarge(self.max_bytes) + return message + + await self.app(scope, limited_receive, send) + + +def _declared_content_length(scope: Scope) -> int | None: + """ + read the Content-Length header from the raw ASGI scope + :param scope: ASGI connection scope + :returns: the declared body length, or None when absent or unparseable + """ + for name, value in scope["headers"]: + if name == b"content-length": + try: + return int(value) + except ValueError: + return None + return None diff --git a/src/hymical_forms/models.py b/src/hymical_forms/models.py new file mode 100644 index 0000000..96284fb --- /dev/null +++ b/src/hymical_forms/models.py @@ -0,0 +1,311 @@ +""" +the persisted schema: endpoints and the submissions addressed to them +""" + +from __future__ import annotations + +from datetime import UTC, datetime + +from sqlalchemy import ( + JSON, + CheckConstraint, + DateTime, + ForeignKey, + String, + TypeDecorator, + UniqueConstraint, +) +from sqlalchemy.engine import Dialect +from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column + +from hymical_forms.ingestion import ( + ENDPOINT_ID_MAX_LENGTH, + IDEMPOTENCY_KEY_MAX_LENGTH, + PAYLOAD_FINGERPRINT_LENGTH, + SUBMISSION_ID_MAX_LENGTH, +) +from hymical_forms.ingestion import Submission as DomainSubmission +from hymical_forms.webhooks import ( + DELIVERY_ATTEMPT_ID_MAX_LENGTH, + DELIVERY_ERROR_MAX_LENGTH, + DELIVERY_STATE_MAX_LENGTH, + WEBHOOK_DELIVERY_ID_MAX_LENGTH, + WEBHOOK_SECRET_MAX_LENGTH, + WEBHOOK_URL_MAX_LENGTH, + DeliveryState, +) + +ENDPOINT_NAME_MAX_LENGTH = 200 +DELIVERY_OUTCOME_MAX_LENGTH = 32 + + +def utcnow() -> datetime: + """ + read the current time as a timezone-aware UTC timestamp + :returns: the current instant in UTC + """ + return datetime.now(UTC) + + +class UtcDateTime(TypeDecorator[datetime]): + """ + a timestamp column that always stores and returns timezone-aware UTC values + """ + + # PostgreSQL round-trips ``TIMESTAMPTZ`` faithfully, but SQLite has no + # timezone-aware storage and hands back a naive datetime. Normalising on the + # way in and out keeps the two backends indistinguishable to the rest of the + # code, so a value written as UTC is always read back as UTC. + impl = DateTime(timezone=True) + cache_ok = True + + def process_bind_param(self, value: datetime | None, dialect: Dialect) -> datetime | None: + """ + convert a value on its way into the database + :param value: the timestamp being stored, or None + :param dialect: the active SQLAlchemy dialect + :returns: the same instant expressed in UTC, or None + :raises ValueError: if the timestamp carries no timezone + """ + if value is None: + return None + if value.tzinfo is None: + raise ValueError("naive datetimes cannot be stored, use a timezone-aware value") + return value.astimezone(UTC) + + def process_result_value(self, value: datetime | None, dialect: Dialect) -> datetime | None: + """ + convert a value on its way out of the database + :param value: the stored timestamp, naive on backends without timezone support + :param dialect: the active SQLAlchemy dialect + :returns: a timezone-aware UTC timestamp, or None + """ + if value is None: + return None + if value.tzinfo is None: + return value.replace(tzinfo=UTC) + return value.astimezone(UTC) + + +class Base(DeclarativeBase): + """ + declarative base for every persisted table + """ + + +class Endpoint(Base): + """ + a form ingestion destination that submissions may be addressed to + """ + + # The public endpoint ID is the primary key. It is already unique, immutable + # in practice (changing it breaks every deployed HTML form pointing at it), + # and constrained to a short safe character set, so a surrogate key would add + # a join without buying anything. + __tablename__ = "endpoints" + + __table_args__ = ( + # An endpoint either has a full webhook configuration or none of it. A URL + # without a secret would mean sending unsigned payloads, which no receiver + # could trust. + CheckConstraint( + "(webhook_url IS NULL) = (webhook_secret IS NULL)", + name="ck_endpoints_webhook_configuration", + ), + ) + + id: Mapped[str] = mapped_column(String(ENDPOINT_ID_MAX_LENGTH), primary_key=True) + name: Mapped[str] = mapped_column(String(ENDPOINT_NAME_MAX_LENGTH)) + is_active: Mapped[bool] = mapped_column(default=True) + created_at: Mapped[datetime] = mapped_column(UtcDateTime, default=utcnow) + + # One destination per endpoint, held here rather than in a table of its own. + # A separate table would only start paying for itself with several + # destinations, which this build deliberately does not have. + webhook_url: Mapped[str | None] = mapped_column(String(WEBHOOK_URL_MAX_LENGTH), default=None) + webhook_secret: Mapped[str | None] = mapped_column( + String(WEBHOOK_SECRET_MAX_LENGTH), default=None + ) + + +class Submission(Base): + """ + a submission that was accepted for a persisted endpoint + """ + + __tablename__ = "submissions" + + __table_args__ = ( + # The database, not the application, is what makes an idempotency key + # unique per endpoint. Two concurrent retries can both find nothing and + # both try to insert, so the constraint is the only authoritative answer. + # + # Both PostgreSQL and SQLite treat NULLs in a unique constraint as + # distinct from each other, so submissions sent without a key stay + # unrestricted without needing a partial index on either backend. + UniqueConstraint( + "endpoint_id", + "idempotency_key", + name="uq_submissions_endpoint_idempotency_key", + ), + # A submission either carries a full idempotency identity or none of it. + CheckConstraint( + "(idempotency_key IS NULL) = (payload_fingerprint IS NULL)", + name="ck_submissions_idempotency_identity", + ), + ) + + id: Mapped[str] = mapped_column(String(SUBMISSION_ID_MAX_LENGTH), primary_key=True) + endpoint_id: Mapped[str] = mapped_column( + String(ENDPOINT_ID_MAX_LENGTH), + ForeignKey("endpoints.id"), + index=True, + ) + received_at: Mapped[datetime] = mapped_column(UtcDateTime) + + # Stored as ``{"field name": ["value", ...]}`` under SQLAlchemy's generic JSON + # type, which is ``json`` on PostgreSQL rather than ``jsonb``. That is + # deliberate: ``jsonb`` normalises object key order, which would silently + # reorder a form's fields, and this payload is written once and read whole, + # so ``jsonb`` indexing would buy nothing here. Every value is a list because + # JSON has no tuple, so the domain's tuples widen here and narrow again in + # :meth:`to_domain`. + fields: Mapped[dict[str, list[str]]] = mapped_column(JSON) + + # Null for submissions sent without an ``Idempotency-Key``, which is the + # common case. When present, the fingerprint is what tells a safe retry apart + # from the same key being reused for different content. + idempotency_key: Mapped[str | None] = mapped_column( + String(IDEMPOTENCY_KEY_MAX_LENGTH), default=None + ) + payload_fingerprint: Mapped[str | None] = mapped_column( + String(PAYLOAD_FINGERPRINT_LENGTH), default=None + ) + + @classmethod + def from_domain( + cls, + submission: DomainSubmission, + *, + idempotency_key: str | None = None, + payload_fingerprint: str | None = None, + ) -> Submission: + """ + build a persistable row from a validated domain submission + :param submission: the normalized submission to store + :param idempotency_key: the client's retry key, or None if it sent none + :param payload_fingerprint: digest of the submitted content, set only alongside a key + :returns: an unsaved row mirroring the submission + """ + return cls( + id=submission.id, + endpoint_id=submission.endpoint_id, + received_at=submission.received_at, + fields={name: list(values) for name, values in submission.fields.items()}, + idempotency_key=idempotency_key, + payload_fingerprint=payload_fingerprint, + ) + + def to_domain(self) -> DomainSubmission: + """ + rebuild the domain submission this row was stored from + :returns: the submission with its repeated field values restored as tuples + """ + return DomainSubmission( + id=self.id, + endpoint_id=self.endpoint_id, + received_at=self.received_at, + fields={name: tuple(values) for name, values in self.fields.items()}, + ) + + +class WebhookDelivery(Base): + """ + the durable obligation to deliver one submission to one destination + """ + + # This is the outbox. It is written in the same transaction as the submission + # it belongs to, so a submission can never be acknowledged while the promise + # to deliver it quietly goes missing. + __tablename__ = "webhook_deliveries" + + __table_args__ = ( + # One logical delivery per submission. An idempotent replay resolves to an + # existing submission, so this constraint is what makes "a replay never + # queues a second delivery" a property of the database rather than of the + # code path that happens to run. + UniqueConstraint("submission_id", name="uq_webhook_deliveries_submission"), + # A delivery is finished exactly when it says it is finished. + CheckConstraint( + "(state IN ('delivered', 'failed')) = (completed_at IS NOT NULL)", + name="ck_webhook_deliveries_completion", + ), + ) + + id: Mapped[str] = mapped_column(String(WEBHOOK_DELIVERY_ID_MAX_LENGTH), primary_key=True) + submission_id: Mapped[str] = mapped_column( + String(SUBMISSION_ID_MAX_LENGTH), ForeignKey("submissions.id") + ) + + # The destination and secret are snapshotted, not read through to the + # endpoint. A delivery represents what was owed when the submission was + # accepted, so changing an endpoint's webhook later must not silently + # redirect work that is already queued, nor leave a queued payload signed + # with a secret its receiver never had. + destination_url: Mapped[str] = mapped_column(String(WEBHOOK_URL_MAX_LENGTH)) + signing_secret: Mapped[str] = mapped_column(String(WEBHOOK_SECRET_MAX_LENGTH)) + + state: Mapped[str] = mapped_column( + String(DELIVERY_STATE_MAX_LENGTH), default=DeliveryState.PENDING + ) + attempts: Mapped[int] = mapped_column(default=0) + + # When this delivery next becomes due. Indexed with the state because that + # pair is exactly what a worker scans for on every poll. + next_attempt_at: Mapped[datetime] = mapped_column(UtcDateTime, index=True) + + # Set while a worker holds the job. A lease that has run out makes the job + # claimable again, which is what stops a worker that died mid-delivery from + # leaving a permanent ``processing`` tombstone. + claim_expires_at: Mapped[datetime | None] = mapped_column(UtcDateTime, default=None) + + created_at: Mapped[datetime] = mapped_column(UtcDateTime, default=utcnow) + completed_at: Mapped[datetime | None] = mapped_column(UtcDateTime, default=None) + + +class DeliveryAttempt(Base): + """ + a record of one outbound request actually made for a delivery + """ + + # The audit trail. One row per request that genuinely went out, so a delivery + # that was retried four times has four rows and none of them is overwritten. + # A job that is merely inspected and found not due produces nothing here. + __tablename__ = "delivery_attempts" + + id: Mapped[str] = mapped_column(String(DELIVERY_ATTEMPT_ID_MAX_LENGTH), primary_key=True) + delivery_id: Mapped[str] = mapped_column( + String(WEBHOOK_DELIVERY_ID_MAX_LENGTH), + ForeignKey("webhook_deliveries.id"), + index=True, + ) + submission_id: Mapped[str] = mapped_column( + String(SUBMISSION_ID_MAX_LENGTH), + ForeignKey("submissions.id"), + index=True, + ) + attempt_number: Mapped[int] = mapped_column() + + # The URL as it was used, not as it is configured now, so the record still + # explains itself after the endpoint's destination changes. + destination_url: Mapped[str] = mapped_column(String(WEBHOOK_URL_MAX_LENGTH)) + attempted_at: Mapped[datetime] = mapped_column(UtcDateTime, default=utcnow) + + # Stored as plain text rather than a database enum, because a database enum + # would need a migration to gain a value and there is no migration tool yet. + outcome: Mapped[str] = mapped_column(String(DELIVERY_OUTCOME_MAX_LENGTH)) + response_status: Mapped[int | None] = mapped_column(default=None) + + # Response bodies are deliberately not stored. They are unbounded, written by + # somebody else's server, and nothing in this build reads them back. + error: Mapped[str | None] = mapped_column(String(DELIVERY_ERROR_MAX_LENGTH), default=None) diff --git a/src/hymical_forms/storage.py b/src/hymical_forms/storage.py new file mode 100644 index 0000000..9b54fe5 --- /dev/null +++ b/src/hymical_forms/storage.py @@ -0,0 +1,392 @@ +""" +persistence operations, the only place queries are written + +Most functions here leave the commit to the caller, so a request handler decides +when its work becomes durable and a failure anywhere before that commit leaves +the database untouched. Three functions own their transaction and say so: +:func:`store_submission`, because it writes a submission and the obligation to +deliver it as one atomic unit and must roll both back together to settle an +idempotency race; :func:`claim_due_deliveries`, because a claim is only worth +anything once it is committed; and :func:`complete_attempt`, because the audit +record and the state it justifies have to land together. +""" + +from __future__ import annotations + +from dataclasses import dataclass +from datetime import datetime, timedelta +from typing import Any, cast + +from sqlalchemy import ColumnElement, and_, or_, select, update +from sqlalchemy.engine import CursorResult +from sqlalchemy.exc import IntegrityError +from sqlalchemy.orm import Session + +from hymical_forms import models +from hymical_forms.ingestion import Submission +from hymical_forms.webhooks import ( + DeliveryOutcome, + DeliveryResult, + DeliveryState, + RetryPolicy, + WebhookTarget, + is_retryable, + new_delivery_attempt_id, + new_webhook_delivery_id, +) + + +class EndpointAlreadyExists(Exception): + """ + raised when an endpoint ID is already taken + """ + + def __init__(self, endpoint_id: str) -> None: + """ + record which endpoint ID was already in use + :param endpoint_id: the identifier that collided + """ + super().__init__(f"endpoint {endpoint_id!r} already exists") + self.endpoint_id = endpoint_id + + +def create_endpoint( + session: Session, + *, + endpoint_id: str, + name: str, + is_active: bool, + webhook_url: str | None = None, + webhook_secret: str | None = None, +) -> models.Endpoint: + """ + add an endpoint, failing if the identifier is taken + :param session: the session to add the endpoint through + :param endpoint_id: the public identifier the endpoint will answer on + :param name: human-readable label for the endpoint + :param is_active: whether the endpoint should accept submissions straight away + :param webhook_url: destination to deliver submissions to, or None for no webhook + :param webhook_secret: signing secret for that destination, set only alongside a URL + :returns: the pending endpoint, not yet committed + :raises EndpointAlreadyExists: if an endpoint already holds that identifier + """ + endpoint = models.Endpoint( + id=endpoint_id, + name=name, + is_active=is_active, + webhook_url=webhook_url, + webhook_secret=webhook_secret, + ) + session.add(endpoint) + try: + # Flushing here turns the unique violation into a catchable error while + # the caller can still react, rather than at an opaque commit later. + session.flush() + except IntegrityError as exc: + session.rollback() + raise EndpointAlreadyExists(endpoint_id) from exc + return endpoint + + +def get_endpoint(session: Session, endpoint_id: str) -> models.Endpoint | None: + """ + look an endpoint up by its public identifier + :param session: the session to query through + :param endpoint_id: the public identifier to resolve + :returns: the endpoint, or None if no endpoint holds that identifier + """ + return session.get(models.Endpoint, endpoint_id) + + +class IdempotencyKeyReused(Exception): + """ + raised when an idempotency key was already used for different content + """ + + def __init__(self, endpoint_id: str, idempotency_key: str) -> None: + """ + record which key was reused, and where + :param endpoint_id: the endpoint the key is scoped to + :param idempotency_key: the key that was already spent on other content + """ + super().__init__(f"idempotency key already used on endpoint {endpoint_id!r}") + self.endpoint_id = endpoint_id + self.idempotency_key = idempotency_key + + +@dataclass(frozen=True, slots=True) +class StoredSubmission: + """ + the outcome of storing a submission + """ + + submission: Submission + replayed: bool + + +def find_by_idempotency_key( + session: Session, endpoint_id: str, idempotency_key: str +) -> models.Submission | None: + """ + look up the submission a key was already spent on + :param session: the session to query through + :param endpoint_id: the endpoint the key is scoped to + :param idempotency_key: the key to resolve + :returns: the earlier submission, or None if the key is unused + """ + return session.scalars( + select(models.Submission).where( + models.Submission.endpoint_id == endpoint_id, + models.Submission.idempotency_key == idempotency_key, + ) + ).one_or_none() + + +def store_submission( + session: Session, + submission: Submission, + *, + now: datetime, + idempotency_key: str | None = None, + payload_fingerprint: str | None = None, + webhook: WebhookTarget | None = None, +) -> StoredSubmission: + """ + store a submission and the obligation to deliver it, or resolve to an earlier one + :param session: the session to write through + :param submission: the validated domain submission to store + :param now: the instant the submission was accepted + :param idempotency_key: the client's retry key, or None if it sent none + :param payload_fingerprint: digest of the submitted content, required alongside a key + :param webhook: the destination owed delivery, or None if the endpoint has no webhook + :returns: the stored submission and whether it came from an earlier attempt + :raises IdempotencyKeyReused: if the key was already used for different content + """ + # The submission and its delivery obligation go in as one transaction. That + # is the whole reliability claim: once this commits, a crash cannot leave a + # submission that was acknowledged with nothing durable saying delivery is + # still owed, and it cannot leave delivery work for a submission that never + # existed either. + if idempotency_key is None: + _add_submission(session, submission, now=now, webhook=webhook) + session.commit() + return StoredSubmission(submission, replayed=False) + + # Fast path for the ordinary retry, where the first attempt already landed. + existing = find_by_idempotency_key(session, submission.endpoint_id, idempotency_key) + if existing is not None: + return StoredSubmission(_settle(existing, payload_fingerprint), replayed=True) + + _add_submission( + session, + submission, + now=now, + webhook=webhook, + idempotency_key=idempotency_key, + payload_fingerprint=payload_fingerprint, + ) + try: + session.commit() + except IntegrityError: + # A concurrent request inserted the same key between the lookup above and + # this commit. The unique constraint is what caught it, which is the whole + # point: the lookup is an optimisation, never the guarantee. + # + # The rollback is mandatory. A session left holding a failed flush refuses + # every later query with PendingRollbackError, so the read below would + # fail rather than find the winner. Rolling back also discards this + # request's submission and its delivery together, leaving exactly the one + # submission and the one delivery the winner committed. + session.rollback() + existing = find_by_idempotency_key(session, submission.endpoint_id, idempotency_key) + if existing is None: + # The violation was something other than the idempotency constraint, + # so it is not ours to interpret and must not be reported as success. + raise + return StoredSubmission(_settle(existing, payload_fingerprint), replayed=True) + + return StoredSubmission(submission, replayed=False) + + +def _add_submission( + session: Session, + submission: Submission, + *, + now: datetime, + webhook: WebhookTarget | None, + idempotency_key: str | None = None, + payload_fingerprint: str | None = None, +) -> None: + """ + stage a submission and, if one is owed, its delivery, without committing + :param session: the session to add through + :param submission: the validated domain submission to store + :param now: the instant the submission was accepted + :param webhook: the destination owed delivery, or None if the endpoint has no webhook + :param idempotency_key: the client's retry key, or None if it sent none + :param payload_fingerprint: digest of the submitted content, set only alongside a key + """ + session.add( + models.Submission.from_domain( + submission, + idempotency_key=idempotency_key, + payload_fingerprint=payload_fingerprint, + ) + ) + if webhook is None: + return + + session.add( + models.WebhookDelivery( + id=new_webhook_delivery_id(), + submission_id=submission.id, + destination_url=webhook.url, + signing_secret=webhook.secret, + state=DeliveryState.PENDING, + attempts=0, + # Due straight away: the first attempt is not delayed, it is simply + # made by a worker rather than by the request that caused it. + next_attempt_at=now, + created_at=now, + ) + ) + + +def due_condition(now: datetime) -> ColumnElement[bool]: + """ + build the test for a delivery a worker is allowed to pick up + :param now: the instant to judge dueness against + :returns: a SQL condition matching claimable deliveries + """ + # Two ways to be claimable: waiting and due, or claimed by a worker whose + # lease has run out. The second is what recovers work from a worker that died + # holding a job, rather than leaving it stuck in ``processing`` forever. + delivery = models.WebhookDelivery + return or_( + and_(delivery.state == DeliveryState.PENDING, delivery.next_attempt_at <= now), + and_(delivery.state == DeliveryState.PROCESSING, delivery.claim_expires_at <= now), + ) + + +def claim_due_deliveries( + session: Session, *, now: datetime, lease_seconds: float, limit: int +) -> list[models.WebhookDelivery]: + """ + take ownership of up to a few due deliveries, in one committed transaction + :param session: the session to claim through + :param now: the instant to judge dueness against + :param lease_seconds: how long the claim protects a delivery from other workers + :param limit: the most deliveries to claim at once + :returns: the deliveries this worker now owns + """ + delivery = models.WebhookDelivery + due = due_condition(now) + + statement = select(delivery).where(due).order_by(delivery.next_attempt_at).limit(limit) + if session.get_bind().dialect.name == "postgresql": + # PostgreSQL can hand each worker a different set of rows outright, which + # is the real answer to two workers scanning at once. SKIP LOCKED means a + # busy row is passed over rather than waited on. + statement = statement.with_for_update(skip_locked=True) + candidates = list(session.scalars(statement)) + + claimed: list[models.WebhookDelivery] = [] + expires_at = now + timedelta(seconds=lease_seconds) + for candidate in candidates: + # The conditional update is the guarantee on backends without row locking: + # whoever gets there first flips the row out of the due condition, and the + # loser's update matches nothing. Redundant under SKIP LOCKED, and cheap. + result = cast( + "CursorResult[Any]", + session.execute( + update(delivery) + .where(delivery.id == candidate.id) + .where(due) + .values(state=DeliveryState.PROCESSING, claim_expires_at=expires_at) + .execution_options(synchronize_session="fetch") + ), + ) + if result.rowcount == 1: + claimed.append(candidate) + + session.commit() + return claimed + + +def load_submissions(session: Session, submission_ids: list[str]) -> dict[str, Submission]: + """ + load the submissions a batch of deliveries is carrying + :param session: the session to query through + :param submission_ids: the submissions to fetch + :returns: the submissions in domain form, keyed by id + """ + # One query for the batch rather than a lookup per delivery. + rows = session.scalars( + select(models.Submission).where(models.Submission.id.in_(submission_ids)) + ) + return {row.id: row.to_domain() for row in rows} + + +def complete_attempt( + session: Session, + delivery: models.WebhookDelivery, + result: DeliveryResult, + *, + now: datetime, + policy: RetryPolicy, +) -> models.DeliveryAttempt: + """ + record one outbound request and move the delivery to whatever it earned + :param session: the session to write through + :param delivery: the delivery the attempt was made for + :param result: what the attempt produced + :param now: the instant the attempt finished + :param policy: how many attempts are allowed and how long to wait between them + :returns: the committed attempt record + """ + # The audit row and the state it justifies are written together, so the + # history can never disagree with the job about how many attempts happened. + attempt_number = delivery.attempts + 1 + attempt = models.DeliveryAttempt( + id=new_delivery_attempt_id(), + delivery_id=delivery.id, + submission_id=delivery.submission_id, + attempt_number=attempt_number, + destination_url=delivery.destination_url, + attempted_at=now, + outcome=str(result.outcome), + response_status=result.response_status, + error=result.error, + ) + session.add(attempt) + + delivery.attempts = attempt_number + delivery.claim_expires_at = None + + if result.outcome is DeliveryOutcome.SUCCEEDED: + delivery.state = DeliveryState.DELIVERED + delivery.completed_at = now + elif is_retryable(result) and not policy.is_exhausted(attempt_number): + delivery.state = DeliveryState.PENDING + delivery.next_attempt_at = now + policy.delay_after(attempt_number) + else: + # Either the receiver said something repeating will not fix, or the + # allowance ran out. Either way this is the last word on the delivery. + delivery.state = DeliveryState.FAILED + delivery.completed_at = now + + session.commit() + return attempt + + +def _settle(existing: models.Submission, payload_fingerprint: str | None) -> Submission: + """ + decide whether an earlier submission is a replay of this one or a clash + :param existing: the submission the key was already spent on + :param payload_fingerprint: digest of the content submitted this time + :returns: the earlier submission, when the content matches + :raises IdempotencyKeyReused: if the content differs from the earlier attempt + """ + if existing.payload_fingerprint != payload_fingerprint: + raise IdempotencyKeyReused(existing.endpoint_id, str(existing.idempotency_key)) + return existing.to_domain() diff --git a/src/hymical_forms/webhooks.py b/src/hymical_forms/webhooks.py new file mode 100644 index 0000000..bf09522 --- /dev/null +++ b/src/hymical_forms/webhooks.py @@ -0,0 +1,306 @@ +""" +webhook rules: destination validation, signing secrets, payload and signature + +Nothing in this module performs I/O. Building and signing a payload is kept +apart from sending it so that the bytes which get signed are provably the bytes +that go on the wire, and so both can be tested without a network. +""" + +from __future__ import annotations + +import hashlib +import hmac +import ipaddress +import json +import secrets +import uuid +from dataclasses import dataclass +from datetime import UTC, datetime, timedelta +from enum import StrEnum +from typing import Any +from urllib.parse import urlsplit + +from hymical_forms.ingestion import Submission + +SIGNATURE_HEADER = "Hymical-Signature" +SIGNATURE_VERSION = "v1" + +SUBMISSION_RECEIVED_EVENT = "submission.received" + +WEBHOOK_SECRET_PREFIX = "whsec_" +WEBHOOK_SECRET_MAX_LENGTH = len(WEBHOOK_SECRET_PREFIX) + 64 + +WEBHOOK_URL_MAX_LENGTH = 2048 +ALLOWED_SCHEMES = ("http", "https") + +DELIVERY_ATTEMPT_ID_PREFIX = "att_" +DELIVERY_ATTEMPT_ID_MAX_LENGTH = len(DELIVERY_ATTEMPT_ID_PREFIX) + 32 + +WEBHOOK_DELIVERY_ID_PREFIX = "whd_" +WEBHOOK_DELIVERY_ID_MAX_LENGTH = len(WEBHOOK_DELIVERY_ID_PREFIX) + 32 +DELIVERY_STATE_MAX_LENGTH = 16 + +# Failure text is written by whatever the destination did, so it is attacker +# influenced and has to be bounded before it reaches a column. +DELIVERY_ERROR_MAX_LENGTH = 500 + + +class DeliveryOutcome(StrEnum): + """ + the coarse result of one webhook delivery attempt + """ + + # These strings are a public contract: they are stored, and reported back on + # the submission response. Exception class names deliberately never appear. + SUCCEEDED = "succeeded" + HTTP_ERROR = "http_error" + TIMEOUT = "timeout" + NETWORK_ERROR = "network_error" + + +class DeliveryState(StrEnum): + """ + where a logical delivery has got to + """ + + # ``pending`` is owed and waiting for its due time, ``processing`` is claimed + # by a worker and holding a lease, and the last two are terminal. + PENDING = "pending" + PROCESSING = "processing" + DELIVERED = "delivered" + FAILED = "failed" + + +TERMINAL_STATES = (DeliveryState.DELIVERED, DeliveryState.FAILED) + +# Statuses below 500 that still describe a condition worth waiting out. Every +# other 4xx is treated as final: repeating a request the receiver called +# malformed, unauthorized or missing will not repair it, and a 409 from a webhook +# receiver almost always means it already has this event. +RETRYABLE_STATUSES = ( + 408, # Request Timeout + 425, # Too Early + 429, # Too Many Requests +) + + +@dataclass(frozen=True, slots=True) +class DeliveryResult: + """ + what one delivery attempt produced + """ + + outcome: DeliveryOutcome + response_status: int | None = None + error: str | None = None + + +def is_retryable(result: DeliveryResult) -> bool: + """ + decide whether an outcome is worth another attempt later + :param result: the outcome of an attempt + :returns: True if the same request might succeed if repeated + """ + if result.outcome is DeliveryOutcome.SUCCEEDED: + return False + if result.outcome in (DeliveryOutcome.TIMEOUT, DeliveryOutcome.NETWORK_ERROR): + return True + + # A 3xx lands here too. Redirects are not followed, and a destination that + # answers with one is misconfigured rather than briefly unwell, so it is + # final like the other non-retryable statuses. + status = result.response_status + return status is not None and (status >= 500 or status in RETRYABLE_STATUSES) + + +@dataclass(frozen=True, slots=True) +class RetryPolicy: + """ + how long to wait between attempts, and when to stop + """ + + max_attempts: int + initial_seconds: float + max_seconds: float + + def delay_after(self, attempts_made: int) -> timedelta: + """ + work out how long to wait before the next attempt + :param attempts_made: how many attempts have already been made + :returns: the wait before the next one becomes due + """ + # Doubling from the initial delay and capped, with no jitter. Jitter would + # spread a thundering herd across a shared destination, but it would also + # make every retry test approximate, and nothing here fans out widely + # enough yet to need it. + exponent = max(attempts_made - 1, 0) + seconds = self.initial_seconds * (2**exponent) + return timedelta(seconds=min(seconds, self.max_seconds)) + + def is_exhausted(self, attempts_made: int) -> bool: + """ + report whether a delivery has used up its allowance + :param attempts_made: how many attempts have already been made + :returns: True if no further attempt should be scheduled + """ + return attempts_made >= self.max_attempts + + +class WebhookUrlRejected(Exception): + """ + raised when a webhook destination is not one this service will send to + """ + + def __init__(self, reason: str) -> None: + """ + record why the destination was refused + :param reason: short phrase completing "the webhook URL ..." + """ + super().__init__(reason) + self.reason = reason + + +def validate_webhook_url(url: str, *, allow_private_targets: bool = False) -> None: + """ + check that a destination is one this service is willing to send to + :param url: the destination the caller wants submissions delivered to + :param allow_private_targets: whether to permit loopback and private addresses + :raises WebhookUrlRejected: if the destination is malformed or not permitted + """ + if len(url) > WEBHOOK_URL_MAX_LENGTH: + raise WebhookUrlRejected(f"must be at most {WEBHOOK_URL_MAX_LENGTH} characters") + + try: + parts = urlsplit(url) + host = parts.hostname + except ValueError as exc: + # urlsplit rejects malformed IPv6 literals and out-of-range ports here. + raise WebhookUrlRejected("must be a well-formed URL") from exc + + if parts.scheme not in ALLOWED_SCHEMES: + raise WebhookUrlRejected("must use the http or https scheme") + if not host: + raise WebhookUrlRejected("must include a host") + + if not allow_private_targets and _is_internal_host(host): + raise WebhookUrlRejected( + "must not address a loopback, private, link-local or otherwise internal host" + ) + + +def _is_internal_host(host: str) -> bool: + """ + report whether a host literal obviously names the server's own network + :param host: the host taken from the destination URL + :returns: True if the host is one submissions must not be delivered to + """ + # Only literals are judged. A name is not resolved here, so a hostname that + # resolves to a private address still passes; see the SSRF note in the README. + name = host.rstrip(".").lower() + if name == "localhost" or name.endswith(".localhost"): + return True + + try: + address = ipaddress.ip_address(name) + except ValueError: + return False + + # ``::ffff:127.0.0.1`` is a loopback address wearing an IPv6 costume, and the + # IPv6 flags do not see through it. + if isinstance(address, ipaddress.IPv6Address) and address.ipv4_mapped is not None: + address = address.ipv4_mapped + + return ( + address.is_loopback + or address.is_private + or address.is_link_local + or address.is_multicast + or address.is_reserved + or address.is_unspecified + ) + + +def new_signing_secret() -> str: + """ + generate a signing secret for a webhook destination + :returns: a prefixed, cryptographically random secret + """ + return f"{WEBHOOK_SECRET_PREFIX}{secrets.token_hex(32)}" + + +def new_delivery_attempt_id() -> str: + """ + generate an opaque identifier for a delivery attempt + :returns: a fresh attempt id such as ``att_1f0c9a...`` + """ + return f"{DELIVERY_ATTEMPT_ID_PREFIX}{uuid.uuid4().hex}" + + +def new_webhook_delivery_id() -> str: + """ + generate an opaque identifier for a logical delivery + :returns: a fresh delivery id such as ``whd_1f0c9a...`` + """ + return f"{WEBHOOK_DELIVERY_ID_PREFIX}{uuid.uuid4().hex}" + + +@dataclass(frozen=True, slots=True) +class WebhookTarget: + """ + the destination and secret a submission is owed delivery to + """ + + # Carried as a pair so that both are snapshotted together when a delivery is + # queued, and neither can be taken from a later version of the endpoint. + url: str + secret: str + + +def build_payload(submission: Submission) -> dict[str, Any]: + """ + build the event body describing a stored submission + :param submission: the submission that was accepted + :returns: the payload to serialize and send + """ + # Repeated values stay lists, exactly as they are stored, so a receiver never + # has to guess whether a field is single or multi valued. + return { + "type": SUBMISSION_RECEIVED_EVENT, + "submission": { + "id": submission.id, + "endpoint_id": submission.endpoint_id, + "received_at": _rfc3339(submission.received_at), + "fields": {name: list(values) for name, values in submission.fields.items()}, + }, + } + + +def serialize_payload(payload: dict[str, Any]) -> bytes: + """ + render a payload to the exact bytes that will be signed and sent + :param payload: the event body to serialize + :returns: the UTF-8 encoded JSON body + """ + return json.dumps(payload, separators=(",", ":"), ensure_ascii=False).encode("utf-8") + + +def sign(body: bytes, secret: str) -> str: + """ + compute the signature header value for an outbound body + :param body: the exact bytes that will be transmitted + :param secret: the destination's signing secret + :returns: the header value, such as ``v1=