Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
60 changes: 57 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ HTTP API providing user/client message handling for an fmsg host. Exposes CRUD o
- [Running](#running)
- [TLS mode (production)](#tls-mode-production)
- [Plain HTTP mode (development / reverse proxy)](#plain-http-mode-development--reverse-proxy)
- [Quotas](#quotas)
- [API Routes](#api-routes)
- [Token and sub-account routes](#post-fmsgtoken)
- [WebSocket route](#get-fmsgws)
Expand Down Expand Up @@ -289,6 +290,38 @@ transfers are not dropped prematurely. These timeouts do not apply to
`/fmsg/ws` connections: once upgraded, a WebSocket connection is hijacked from
the HTTP server and kept alive by its own ping/pong heartbeat.

## Quotas

fmsgid holds each address's limits and usage (`-1` is unlimited; an address in
a quota pool has its pool owner's limits and the whole pool's usage). fmsgd
checks and records mail it receives from other hosts. This service checks and
records the rest: the mail it sends, and the mail it delivers between local
addresses, which never passes through fmsgd.

- **Sending.** Before a message is sent (`POST /fmsg/:id/send`, and the
reaction message `POST /fmsg/:id/react` sends), the sender's current limits
and usage are fetched from fmsgid and the message refused if it is bigger
than the per-message send size (`413`), or would exceed the messages sent per
day or bytes sent per day (`429`) or total bytes sent (`422`). Once sent, the
message is recorded once with `POST /fmsgid/send`, however many recipients it
has. Adding recipients with add-to doesn't send a new message, so it isn't
counted against send limits; each new local recipient is checked and counted
when it receives the message.
- **Local delivery.** Each local recipient gets the checks fmsgd makes: unknown
(`100`), not accepting new messages (`102`), or over its messages received per
day, bytes received per day or total bytes received for one more message of
this size (`101`, user full). Those recipients aren't delivered, and their
code is recorded as their `response_code` in `to_delivery`/`add_to`. Each
recipient that is delivered is recorded with `POST /fmsgid/recv`.

Sizes are the stored size: the body plus every attachment. Usage is
timestamped when this service sent or delivered the message. The checks use
fmsgid's current usage, never the lookup cached for authentication, which only
holds whether the address exists and accepts new messages. If fmsgid can't be
reached a send is refused with `503`, and a local recipient is left pending.
Failing to record usage is logged and never fails the request, since the
message has already been sent.

## API Routes

All routes are prefixed with `/fmsg`. `POST /fmsg/token` accepts an API key and
Expand Down Expand Up @@ -755,17 +788,29 @@ third-party delivery's outcome, so the reply is attempted and the wire's
originating domain always passes (it retains its outgoing messages). Local
recipients are unaffected.

The sender's send limits are then checked against the message's size, the
body plus its attachments (see [Quotas](#quotas)). A refusal leaves the
message a draft, and carries a `code` naming the limit:
`{"error": "daily message limit reached", "code": "send_count_per_1d_limit"}`.

Local recipients are delivered at once, each subject to its receive limits;
a recipient over one is not delivered and its `response_code` is `101`.

**Response:** `200 OK` with `{"id": <int>, "time": <float64>, "sha256": "<64 hex characters>"}`.

**Errors:**

| Status | Condition |
| ------ | --------- |
| `400` | Message is not sendable (no recipients, invalid recipient address, no type, unsupported version) |
| `403` | Not the owner |
| `403` | Not the owner, or the sender is no longer in fmsgid or accepting new messages |
| `404` | Message not found |
| `409` | Message already sent |
| `409` | Reply can never be accepted by one or more remote recipient domains (parent never addressed there, or its delivery there failed) |
| `413` | Message is bigger than the sender's per-message send size (`send_size_per_msg_limit`) |
| `422` | Sender's total sent bytes would exceed its limit, `storage limit reached` (`send_size_total_limit`) |
| `429` | Sender's messages per day (`send_count_per_1d_limit`) or bytes per day (`send_size_per_1d_limit`) would exceed its limit |
| `503` | fmsgid unavailable, so the send limits can't be checked |

### POST `/fmsg/:id/read`

Expand Down Expand Up @@ -795,6 +840,11 @@ This endpoint records the add-to as a new `msg_add_to_batch` row (capturing the

New addresses must be distinct among themselves (case-insensitive).

On a sent message, new local recipients are delivered at once, each subject to
its receive limits (a recipient over one gets `response_code` `101`). Adding
recipients isn't counted against the caller's send limits (see
[Quotas](#quotas)).

**Response:** `200 OK` with `{"id": <int>, "added": <int>, "batch_id": <int>, "sha256": "<64 hex characters>"}`. On a draft, the batch hash is `null` until send finalizes it.

**Errors:**
Expand All @@ -816,8 +866,10 @@ in `to` or `add_to` — and the message must be sent and not terminal.
The reaction is sent immediately as a reaction message: a reply to `:id` with
`no_reply` and `terminal` set, type `text/plain;charset=UTF-8`, body the emoji,
addressed to every other participant of the message. Delivery then proceeds
exactly as for any sent message. A reaction message cannot be replied to,
reacted to, or have recipients added.
exactly as for any sent message, and like any message it counts against the
reactor's send limits and each local recipient's receive limits (see
[Quotas](#quotas)). A reaction message cannot be replied to, reacted to, or
have recipients added.

**Request body (JSON):**

Expand All @@ -839,6 +891,8 @@ the caller has no reaction returns `200 OK` with `null` values.
| `403` | Authenticated identity is not a participant |
| `404` | Message not found |
| `409` | Message is a draft, is terminal, has no other participants, or a remote recipient host cannot accept the reaction because it does not hold the message |
| `413`, `422`, `429` | The reaction would exceed one of the reactor's send limits, as for send |
| `503` | fmsgid unavailable, so the send limits can't be checked |

### GET `/fmsg/:id/data`

Expand Down
14 changes: 10 additions & 4 deletions internal/handlers/finalization_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,15 @@ type finalizationAPI struct {
}

func newFinalizationAPI(t *testing.T) *finalizationAPI {
t.Helper()
return newFinalizationAPIWithID(t, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
fmt.Fprint(w, `{"acceptingNew":true}`)
}))
}

// newFinalizationAPIWithID is newFinalizationAPI with fmsgid served by fmsgID.
func newFinalizationAPIWithID(t *testing.T, fmsgID http.Handler) *finalizationAPI {
t.Helper()
dsn := os.Getenv("FMSG_TEST_DATABASE_URL")
if dsn == "" {
Expand Down Expand Up @@ -66,10 +75,7 @@ func newFinalizationAPI(t *testing.T) *finalizationAPI {
if _, err = pool.Exec(ctx, string(sql)); err != nil {
t.Fatal(err)
}
id := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
fmt.Fprint(w, `{"acceptingNew":true}`)
}))
id := httptest.NewServer(fmsgID)
t.Cleanup(id.Close)
h := NewMessageHandler(&db.DB{Pool: pool}, t.TempDir(), 1<<20, 2<<20, 256, nil, id.URL, "example.com")
a := NewAttachmentHandler(h.DB, h.DataDir, 1<<20, 2<<20)
Expand Down
65 changes: 47 additions & 18 deletions internal/handlers/messages.go
Original file line number Diff line number Diff line change
Expand Up @@ -112,40 +112,56 @@ func formatDeliveredISO(t *float64) *string {
// skips the local domain (no network hop is needed), so it never sets
// time_delivered/response_code for same-domain recipients — fmsg-webapi has
// to resolve that status itself via fmsgid, the same source fmsgd would have
// consulted. Errors talking to fmsgid are logged and left pending (no status
// written) rather than treated as a delivery failure.
// consulted, with the same checks: unknown (100), not accepting (102), or
// over a receive limit (101, user full) for one more message of the message's
// stored size. A delivered recipient's usage is recorded with fmsgid, at the
// delivery time. fmsgid is asked afresh for every recipient (not the cached
// lookup), since usage changes with every message. Errors talking to fmsgid
// are logged and left pending (no status written) rather than treated as a
// delivery failure; a failure to record usage is only logged.
func (h *MessageHandler) resolveLocalDelivery(ctx context.Context, table string, msgID int64, localDomain string, addrs []string) {
now := float64(time.Now().UnixMicro()) / 1e6
query := "UPDATE " + table + " SET time_delivered = $1, response_code = $2 WHERE msg_id = $3 AND addr = $4 AND time_delivered IS NULL AND response_code IS NULL"
size := int64(-1) // the message's stored size, loaded once it is needed
for _, addr := range addrs {
_, domain := parseAddr(addr)
if !strings.EqualFold(domain, localDomain) {
continue // remote — fmsgd handles this
}
code, accepting, err := middleware.CheckFmsgID(h.IDURL, addr)
detail, status, err := middleware.FetchFmsgIDDetail(h.IDURL, addr)
if err != nil {
log.Printf("resolve local delivery: fmsgid check %s: %v", addr, err)
continue
}
if size < 0 && status == http.StatusOK {
if size, err = h.messageStoredSize(ctx, msgID); err != nil {
log.Printf("resolve local delivery: size of msg %d: %v", msgID, err)
size = -1
continue
}
}
code, ok := localDeliveryCode(status, detail, size)
if !ok {
log.Printf("resolve local delivery: unexpected fmsgid status %d for %s", status, addr)
continue
}
var delivered *float64
var responseCode *int
switch {
case code == http.StatusNotFound:
rc := 100 // user unknown
responseCode = &rc
case code == http.StatusOK && !accepting:
rc := 102 // user not accepting
responseCode = &rc
case code == http.StatusOK && accepting:
rc := 200 // accept
responseCode = &rc
if code == responseCodeAccept {
delivered = &now
default:
log.Printf("resolve local delivery: unexpected fmsgid status %d for %s", code, addr)
continue
} else if code == responseCodeUserFull {
log.Printf("resolve local delivery: msg %d to %s refused: over a receive limit", msgID, addr)
}
if _, err := h.q().Exec(ctx, query, delivered, responseCode, msgID, addr); err != nil {
tag, err := h.q().Exec(ctx, query, delivered, code, msgID, addr)
if err != nil {
log.Printf("resolve local delivery: update %s msg %d addr %s: %v", table, msgID, addr, err)
continue
}
// Record usage only for the delivery made here, never for a row
// another request had already resolved.
if delivered != nil && tag.RowsAffected() > 0 {
if err := middleware.RecordFmsgIDRecv(h.IDURL, addr, now, size); err != nil {
log.Printf("resolve local delivery: record recv usage %s msg %d: %v", addr, msgID, err)
}
}
}
}
Expand Down Expand Up @@ -961,6 +977,18 @@ func (h *MessageHandler) Send(c *gin.Context) {
}
}

// The sender's send limits, for the message's stored size (body plus
// attachments). The message is counted once, however many recipients.
size, err := h.messageStoredSize(ctx, msgID)
if err != nil {
log.Printf("send message %d: size: %v", msgID, err)
c.JSON(http.StatusInternalServerError, gin.H{"error": "failed to retrieve message"})
return
}
if !h.checkSendQuota(c, existing.From, size) {
return
}

now := float64(time.Now().UnixMicro()) / 1e6
tx, ok := h.query.(pgx.Tx)
if !ok {
Expand All @@ -981,6 +1009,7 @@ func (h *MessageHandler) Send(c *gin.Context) {
afterCommit(c, func() {
plain := *h
plain.query = nil
plain.recordSend(existing.From, now, size)
plain.resolveLocalDelivery(ctx, "msg_to", msgID, h.LocalDomain, existing.To)
for _, b := range existing.AddTo {
plain.resolveLocalDelivery(ctx, "msg_add_to", msgID, h.LocalDomain, b.To)
Expand Down
141 changes: 141 additions & 0 deletions internal/handlers/quota.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,141 @@
package handlers

import (
"context"
"fmt"
"log"
"net/http"

"github.com/gin-gonic/gin"
"github.com/markmnl/fmsg-webapi/internal/middleware"
)

// Quotas. fmsgid holds each address's limits and usage; this service checks
// and records the messages that don't pass through fmsgd's receive path:
//
// - Sending (send and react): the sender's send limits are checked before
// the message is sent, and the message is recorded once with POST
// /fmsgid/send after it is sent, whatever its number of recipients.
// - Local delivery (see resolveLocalDelivery): each local recipient's
// receive limits are checked as fmsgd checks them, a recipient over a
// limit is recorded with response code 101 (user full) and not delivered,
// and each delivered recipient is recorded with POST /fmsgid/recv.
//
// Every size is the stored size: the message body plus its attachments.
// Usage is timestamped when this service sent or delivered the message.
// Recording failures are logged and never fail the request, because the
// message has already been sent.

// Per-recipient response codes (fmsg specification) that local delivery
// records in msg_to/msg_add_to.response_code, as fmsgd does for recipients of
// the messages it receives.
const (
responseCodeUserUnknown = 100
responseCodeUserFull = 101
responseCodeUserNotAccepting = 102
responseCodeAccept = 200
)

// recvLimitExceeded reports whether one more message of size bytes would take
// an address over any of its receive limits (-1 is unlimited). These are the
// checks fmsgd makes for a recipient before accepting a message.
func recvLimitExceeded(d *middleware.FmsgIDDetail, size int64) bool {
return (d.LimitRecvCountPer1d > -1 && d.RecvCountPer1d+1 > d.LimitRecvCountPer1d) ||
(d.LimitRecvSizePer1d > -1 && d.RecvSizePer1d+size > d.LimitRecvSizePer1d) ||
(d.LimitRecvSizeTotal > -1 && d.RecvSizeTotal+size > d.LimitRecvSizeTotal)
}

// localDeliveryCode is the response code local delivery records for a
// recipient, given fmsgid's status and detail for it and the message's stored
// size. ok is false for an unexpected fmsgid status, which leaves the
// recipient pending.
func localDeliveryCode(status int, d *middleware.FmsgIDDetail, size int64) (code int, ok bool) {
switch {
case status == http.StatusNotFound:
return responseCodeUserUnknown, true
case status != http.StatusOK || d == nil:
return 0, false
case !d.AcceptingNew:
return responseCodeUserNotAccepting, true
case recvLimitExceeded(d, size):
return responseCodeUserFull, true
default:
return responseCodeAccept, true
}
}

// sendLimitError is the refusal for a message over one of its sender's send
// limits.
type sendLimitError struct {
status int
code string
message string
}

// checkSendLimits checks one more message of size bytes against a sender's
// send limits (-1 is unlimited), and returns the refusal for the first limit
// it would exceed, or nil. A message too big to ever send is 413, a daily
// limit is 429 (it frees up as the day's messages age out), and the sent-mail
// storage limit is 422.
func checkSendLimits(d *middleware.FmsgIDDetail, size int64) *sendLimitError {
switch {
case d.LimitSendSizePerMsg > -1 && size > d.LimitSendSizePerMsg:
return &sendLimitError{http.StatusRequestEntityTooLarge, "send_size_per_msg_limit",
fmt.Sprintf("message size %d exceeds the sending limit of %d bytes per message", size, d.LimitSendSizePerMsg)}
case d.LimitSendCountPer1d > -1 && d.SendCountPer1d+1 > d.LimitSendCountPer1d:
return &sendLimitError{http.StatusTooManyRequests, "send_count_per_1d_limit", "daily message limit reached"}
case d.LimitSendSizePer1d > -1 && d.SendSizePer1d+size > d.LimitSendSizePer1d:
return &sendLimitError{http.StatusTooManyRequests, "send_size_per_1d_limit", "daily sending size limit reached"}
case d.LimitSendSizeTotal > -1 && d.SendSizeTotal+size > d.LimitSendSizeTotal:
return &sendLimitError{http.StatusUnprocessableEntity, "send_size_total_limit", "storage limit reached"}
}
return nil
}

// checkSendQuota checks that sender may send one more message of size bytes,
// with fmsgid's current usage (never the cached lookup). Otherwise it answers
// with the refusal and returns false.
func (h *MessageHandler) checkSendQuota(c *gin.Context, sender string, size int64) bool {
detail, status, err := middleware.FetchFmsgIDDetail(h.IDURL, sender)
if err != nil {
log.Printf("send quota: fmsgid detail %s: %v", sender, err)
c.JSON(http.StatusServiceUnavailable, gin.H{"error": "identity service unavailable"})
return false
}
switch {
case status == http.StatusNotFound:
c.JSON(http.StatusForbidden, gin.H{"error": fmt.Sprintf("User %s not found", sender)})
return false
case status != http.StatusOK:
log.Printf("send quota: unexpected fmsgid status %d for %s", status, sender)
c.JSON(http.StatusServiceUnavailable, gin.H{"error": "identity service unavailable"})
return false
case !detail.AcceptingNew:
c.JSON(http.StatusForbidden, gin.H{"error": fmt.Sprintf("User %s not authorised to send new messages", sender)})
return false
}
if refusal := checkSendLimits(detail, size); refusal != nil {
log.Printf("send quota: %s refused: %s", sender, refusal.code)
c.JSON(refusal.status, gin.H{"error": refusal.message, "code": refusal.code})
return false
}
return true
}

// recordSend records one sent message of size bytes for sender at ts. A
// failure is logged: the message has been sent.
func (h *MessageHandler) recordSend(sender string, ts float64, size int64) {
if err := middleware.RecordFmsgIDSend(h.IDURL, sender, ts, size); err != nil {
log.Printf("record send usage: %s: %v", sender, err)
}
}

// messageStoredSize is a message's stored size: its body plus its
// attachments.
func (h *MessageHandler) messageStoredSize(ctx context.Context, msgID int64) (int64, error) {
var size int64
err := h.q().QueryRow(ctx,
`SELECT m.size + COALESCE((SELECT sum(a.filesize) FROM msg_attachment a WHERE a.msg_id = m.id), 0)
FROM msg m WHERE m.id = $1`, msgID).Scan(&size)
return size, err
}
Loading
Loading