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
7 changes: 7 additions & 0 deletions docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -143,6 +143,13 @@ reflection boundary; generators do not reach into runtime structs. This keeps
code generation deterministic and lets the server's concurrency machinery
evolve without turning internal representation into an accidental API.

The terminal-ownership state machine for active requests lives in the private
`backend-pending` module. It owns admission capacity, exactly-once terminal
claims, cancellation, and the short State-commit barrier. The transport server
owns frames, custodians, and the bounded writer queue. Direct characterization
tests cover the private state machine so lifecycle changes do not require
reaching through the public API or growing `backend.rkt` further.

Before 1.0, callers that used `registered-*` or `*-info-*` should migrate to
the corresponding `backend-schema` entry. Application State access should use
`state-ref` and `state-set!`; `state?` is available when a predicate is needed.
Expand Down
110 changes: 10 additions & 100 deletions rivet/backend.rkt
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
racket/async-channel
racket/list
racket/match
"private/backend-pending.rkt"
"private/event-id.rkt"
"protocol.rkt")

Expand Down Expand Up @@ -39,7 +40,6 @@
;; details; applications only need a predicate plus state-ref/state-set!.
(define (state? value) (state-info? value))

(struct pending-request (custodian terminal-owned cancel-deferred cancel-requested) #:mutable)
(struct record-info (name field-names field-types) #:transparent)
(struct record-value (name fields) #:transparent)
(struct enum-info (name cases) #:transparent)
Expand Down Expand Up @@ -836,8 +836,15 @@
(define runtime-custodian (make-custodian server-custodian))
(define writer-custodian (make-custodian server-custodian))
(define responses (make-async-channel max-outgoing-frames))
(define pending (make-hash))
(define pending-lock (make-semaphore 1))
(define pending (make-pending-table max-pending-requests))
(define request-id-pending? (pending-table-id-pending? pending))
(define admit-request! (pending-table-admit! pending))
(define claim-pending! (pending-table-claim! pending))
(define cancel-action! (pending-table-cancel-action! pending))
(define begin-state-commit! (pending-table-begin-state-commit! pending))
(define end-state-commit! (pending-table-end-state-commit! pending))
(define release-pending! (pending-table-release! pending))
(define take-all-pending! (pending-table-take-all! pending))
(define event-id-lock (make-semaphore 1))
(define writer-error (box #f))
(define reader-error (box #f))
Expand Down Expand Up @@ -882,103 +889,6 @@
(raise-writer-failure!))
(void))

(define (request-id-pending? id)
(call-with-semaphore
pending-lock
(lambda () (hash-has-key? pending id))))

(define (admit-request! id custodian)
(call-with-semaphore
pending-lock
(lambda ()
(cond
[(hash-has-key? pending id) 'duplicate]
[(>= (hash-count pending) max-pending-requests) 'full]
[else
(hash-set! pending id (pending-request custodian #f #f #f))
'admitted]))))

(define (claim-pending! id)
;; Completion/error/cancellation claim terminal ownership without removing
;; the entry yet. The request continues to occupy its pending slot while a
;; terminal frame is waiting for output capacity, so backpressure cannot be
;; bypassed by admitting an unbounded stream of newly completed requests.
(call-with-semaphore
pending-lock
(lambda ()
(define request (hash-ref pending id #f))
(cond
[(and request
(not (pending-request-terminal-owned request)))
(set-pending-request-terminal-owned! request #t)
request]
[else #f]))))

(define (cancel-action! id)
;; State commits can briefly defer cancellation after the cell becomes
;; visible and until its reserved Event is accepted by the output queue.
;; Keep the pending slot occupied during that interval so cancellation
;; cannot bypass either correlation ownership or the concurrency limit.
(call-with-semaphore
pending-lock
(lambda ()
(define request (hash-ref pending id #f))
(cond
[(or (not request)
(pending-request-terminal-owned request))
#f]
[(pending-request-cancel-deferred request)
(set-pending-request-cancel-requested! request #t)
'deferred]
[else
(set-pending-request-terminal-owned! request #t)
request]))))

(define (begin-state-commit! id)
(call-with-semaphore
pending-lock
(lambda ()
(define request (hash-ref pending id #f))
(cond
[(not request) 'untracked]
[(pending-request-terminal-owned request) 'terminal]
[else
(set-pending-request-cancel-deferred! request #t)
request]))))

(define (end-state-commit! id request)
;; Clear the barrier and atomically convert any deferred Cancel into terminal
;; ownership before another Cancel or normal Response can race in.
(call-with-semaphore
pending-lock
(lambda ()
(define current (hash-ref pending id #f))
(cond
[(not (eq? current request)) #f]
[else
(set-pending-request-cancel-deferred! request #f)
(cond
[(and (pending-request-cancel-requested request)
(not (pending-request-terminal-owned request)))
(set-pending-request-terminal-owned! request #t)
request]
[else #f])]))))

(define (release-pending! id request)
(call-with-semaphore
pending-lock
(lambda ()
(when (eq? (hash-ref pending id #f) request)
(hash-remove! pending id)))))

(define (take-all-pending!)
(call-with-semaphore
pending-lock
(lambda ()
(define requests (hash-values pending))
(hash-clear! pending)
(map pending-request-custodian requests))))

(define (send-claimed! id request response)
;; Always release the pending slot after the send attempt, including an
;; asynchronous break or writer failure. On transport failure the outer
Expand Down
136 changes: 136 additions & 0 deletions rivet/private/backend-pending.rkt
Original file line number Diff line number Diff line change
@@ -0,0 +1,136 @@
#lang racket/base

(provide make-pending-table
(struct-out pending-table)
pending-request?
pending-request-custodian)

;; Pending requests have one terminal owner. Completion, failure, and
;; cancellation race through this table, while a short state-commit barrier can
;; defer cancellation until the matching State event has entered the output
;; queue. Keeping this state machine separate makes its invariants directly
;; testable without exposing it from rivet/backend.
(struct pending-request
(custodian terminal-owned cancel-deferred cancel-requested)
#:mutable)

(struct pending-table
(id-pending?
admit!
claim!
cancel-action!
begin-state-commit!
end-state-commit!
release!
take-all!)
#:transparent)

(define (make-pending-table limit)
(unless (exact-positive-integer? limit)
(raise-argument-error 'make-pending-table "exact-positive-integer?" limit))

(define entries (make-hash))
(define lock (make-semaphore 1))

(define (id-pending? id)
(call-with-semaphore
lock
(lambda () (hash-has-key? entries id))))

(define (admit! id custodian)
(call-with-semaphore
lock
(lambda ()
(cond
[(hash-has-key? entries id) 'duplicate]
[(>= (hash-count entries) limit) 'full]
[else
(hash-set! entries id (pending-request custodian #f #f #f))
'admitted]))))

(define (claim! id)
;; Claim terminal ownership without freeing capacity. A completed request
;; continues to occupy its slot until its terminal frame is accepted by the
;; bounded output queue, so backpressure cannot be bypassed.
(call-with-semaphore
lock
(lambda ()
(define request (hash-ref entries id #f))
(cond
[(and request
(not (pending-request-terminal-owned request)))
(set-pending-request-terminal-owned! request #t)
request]
[else #f]))))

(define (cancel-action! id)
(call-with-semaphore
lock
(lambda ()
(define request (hash-ref entries id #f))
(cond
[(or (not request)
(pending-request-terminal-owned request))
#f]
[(pending-request-cancel-deferred request)
(set-pending-request-cancel-requested! request #t)
'deferred]
[else
(set-pending-request-terminal-owned! request #t)
request]))))

(define (begin-state-commit! id)
(call-with-semaphore
lock
(lambda ()
(define request (hash-ref entries id #f))
(cond
[(not request) 'untracked]
[(pending-request-terminal-owned request) 'terminal]
[else
(set-pending-request-cancel-deferred! request #t)
request]))))

(define (end-state-commit! id request)
;; Clear the barrier and atomically convert a deferred Cancel into terminal
;; ownership before another Cancel or normal Response can race in.
(call-with-semaphore
lock
(lambda ()
(define current (hash-ref entries id #f))
(cond
[(not (eq? current request)) #f]
[else
(set-pending-request-cancel-deferred! request #f)
(cond
[(and (pending-request-cancel-requested request)
(not (pending-request-terminal-owned request)))
(set-pending-request-terminal-owned! request #t)
request]
[else #f])]))))

(define (release! id request)
;; Identity guards against an old worker releasing a newer request that
;; reused the same wire correlation id.
(call-with-semaphore
lock
(lambda ()
(when (eq? (hash-ref entries id #f) request)
(hash-remove! entries id)))))

(define (take-all!)
(call-with-semaphore
lock
(lambda ()
(define requests (hash-values entries))
(hash-clear! entries)
(map pending-request-custodian requests))))

(pending-table id-pending?
admit!
claim!
cancel-action!
begin-state-commit!
end-state-commit!
release!
take-all!))
100 changes: 100 additions & 0 deletions tests/backend-pending.rkt
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
#lang racket/base

(require rackunit
"../rivet/private/backend-pending.rkt")

(define (operations table)
(values (pending-table-id-pending? table)
(pending-table-admit! table)
(pending-table-claim! table)
(pending-table-cancel-action! table)
(pending-table-begin-state-commit! table)
(pending-table-end-state-commit! table)
(pending-table-release! table)
(pending-table-take-all! table)))

(check-exn exn:fail:contract? (lambda () (make-pending-table 0)))

;; Admission is bounded and duplicate identifiers never replace their owner.
(let ()
(define table (make-pending-table 1))
(define-values (pending? admit! claim! cancel! begin! end! release! take-all!)
(operations table))
(define first-owner (make-custodian))
(define other-owner (make-custodian))
(check-eq? (admit! 7 first-owner) 'admitted)
(check-eq? (admit! 7 other-owner) 'duplicate)
(check-eq? (admit! 8 other-owner) 'full)
(check-true (pending? 7))

;; A claim is exactly once and retains capacity until the sender releases it.
(define request (claim! 7))
(check-true (pending-request? request))
(check-eq? (pending-request-custodian request) first-owner)
(check-false (claim! 7))
(check-false (cancel! 7))
(check-eq? (admit! 8 other-owner) 'full)
(release! 7 request)
(check-false (pending? 7))
(check-eq? (admit! 8 other-owner) 'admitted)
(for ([owner (in-list (take-all!))])
(custodian-shutdown-all owner)))

;; Cancellation can own the terminal outcome before normal completion.
(let ()
(define table (make-pending-table 2))
(define-values (pending? admit! claim! cancel! begin! end! release! take-all!)
(operations table))
(define owner (make-custodian))
(check-eq? (admit! 11 owner) 'admitted)
(define cancelled (cancel! 11))
(check-true (pending-request? cancelled))
(check-false (cancel! 11))
(check-false (claim! 11))
(release! 11 cancelled)
(check-false (pending? 11))
(custodian-shutdown-all owner))

;; A State commit defers cancellation until the event side effect is admitted,
;; then hands terminal ownership to exactly one cancellation path.
(let ()
(define table (make-pending-table 1))
(define-values (pending? admit! claim! cancel! begin! end! release! take-all!)
(operations table))
(define owner (make-custodian))
(check-eq? (admit! 19 owner) 'admitted)
(define request (begin! 19))
(check-true (pending-request? request))
(check-eq? (cancel! 19) 'deferred)
(check-eq? (cancel! 19) 'deferred)
(check-eq? (end! 19 request) request)
(check-false (end! 19 request))
(check-false (claim! 19))

;; Releasing with the wrong identity cannot remove the current request.
(define wrong-table (make-pending-table 1))
(define wrong-admit! (pending-table-admit! wrong-table))
(define wrong-claim! (pending-table-claim! wrong-table))
(check-eq? (wrong-admit! 19 (make-custodian)) 'admitted)
(release! 19 (wrong-claim! 19))
(check-true (pending? 19))
(release! 19 request)
(check-false (pending? 19))
(custodian-shutdown-all owner)
(for ([remaining-owner (in-list ((pending-table-take-all! wrong-table)))])
(custodian-shutdown-all remaining-owner)))

;; Shutdown drains every custodian and atomically clears correlation ownership.
(let ()
(define table (make-pending-table 3))
(define admit! (pending-table-admit! table))
(define pending? (pending-table-id-pending? table))
(define take-all! (pending-table-take-all! table))
(define owners (for/list ([id '(1 2 3)]) (make-custodian)))
(for ([id '(1 2 3)] [owner (in-list owners)])
(check-eq? (admit! id owner) 'admitted))
(check-equal? (sort (take-all!) < #:key eq-hash-code)
(sort owners < #:key eq-hash-code))
(for ([id '(1 2 3)]) (check-false (pending? id)))
(check-equal? (take-all!) '())
(for ([owner (in-list owners)]) (custodian-shutdown-all owner)))
Loading