fix: persist dead-letter payload, periodic pruning, atomic replay claim, remove dead code, 503 quota
ober
54f3af19a4752a31bd247c15e8e9553cac3205ee
--- a/edge.ss +++ b/edge.ss @@ -143,6 +143,7 @@ (def *replay-window-seconds* (env-positive-integer "EDGE_REPLAY_WINDOW_SECONDS" 300 30 3600)) (def *event-ttl-seconds* (env-positive-integer "EDGE_EVENT_TTL_SECONDS" 86400 60 31536000)) (def *dead-ttl-seconds* (env-positive-integer "EDGE_DEAD_LETTER_TTL_SECONDS" 604800 60 31536000)) +(def *prune-interval-seconds* (env-positive-integer "EDGE_PRUNE_INTERVAL_SECONDS" 60 1 3600)) (def *max-requests-per-second* (env-positive-integer "EDGE_MAX_REQUESTS_PER_SECOND" 200 1 100000)) (def (loopback-host? host) @@ -339,6 +340,15 @@ (when claimed? (db-note-write-and-prune!)) claimed?))))) +;; Roll back a replay-key claim. Used when a two-key claim is only partially +;; successful so the restart-resistant DB never disagrees with the in-memory +;; cache about whether an event identity has been seen. +(def (db-release-replay! replay-key) + (when *db* + (with-mutex *db-mutex* + (guard (e [(condition? e) (void)]) + (sqlite-execute *db* "DELETE FROM replay_keys WHERE replay_key = ?" replay-key))))) + (def (db-load-events!) ;; Restore only the bounded, unexpired retention window into memory. (when *db* @@ -393,18 +403,6 @@ (db-close-handle! db)))))) ;; ═══════════════════════════════════════════════════════════════ -;; ID Generation -;; ═══════════════════════════════════════════════════════════════ - -(def id-counter 0) -(def id-mutex (make-mutex)) - -(def (gen-id) - (with-mutex id-mutex - (set! id-counter (+ id-counter 1)) - (format "evt_~a_~a" (time-second (current-time)) id-counter))) - -;; ═══════════════════════════════════════════════════════════════ ;; State Store (STM) ;; ;; All application state lives in STM refs. Every mutation is @@ -498,6 +496,31 @@ (ref-set events-ref events) (ref-set dead-letter-ref dead)))))) +;; Pruning is an O(events) scan under the global state mutex, so it runs on a +;; periodic timer rather than inline on every read path (get-event, event-stats, +;; get-dead-letter-list). Reads stay lock-light; expiry is eventually consistent +;; within *prune-interval-seconds*. +(def *prune-running?* #f) + +(def (run-prune-scheduler-guarded) + (let loop () + (when *prune-running?* + (guard (exn [(condition? exn) + (jlog "prune scheduler iteration failed" "err" + (if (message-condition? exn) + (condition-message exn) (format "~a" exn)))]) + (prune-expired-state!)) + (sleep (make-time 'time-duration 0 *prune-interval-seconds*)) + (loop)))) + +(def (start-prune-scheduler!) + (unless *prune-running?* + (set! *prune-running?* #t) + (fork-thread run-prune-scheduler-guarded))) + +(def (stop-prune-scheduler!) + (set! *prune-running?* #f)) + (def (store-result! event status) (let* ([id (hash-ref event "id")] [record @@ -510,13 +533,11 @@ (db-store-event! event status))) (def (get-event id) - (prune-expired-state!) (let ([st (ref-deref events-ref)]) (and (persistent-map-has? st id) (persistent-map-ref st id)))) (def (event-stats) - (prune-expired-state!) (let ([st (ref-deref events-ref)] [total 0] [ok-n 0] [err-n 0]) (persistent-map-for-each @@ -535,18 +556,20 @@ (def (store-dead-letter! event reason) (let* ([id (hash-ref event "id")] [safe-reason (bounded-string reason 1024)] + [payload (hash-ref event "payload" #f)] + [payload-json (if (hash-table? payload) (json-object->string payload) "{}")] [record (list->hash-table `(("id" . ,id) ("type" . ,(hash-ref event "type" "unknown")) ("reason" . ,safe-reason) ("retries" . ,(hash-ref event "retry-count" 0)) + ("payload" . ,payload-json) ("time" . ,(format "~a" (time-second (current-time))))))]) (with-mutex *state-mutex* (%bounded-dead-put! id record)) (db-store-dead-letter! event safe-reason))) (def (get-dead-letter-list) - (prune-expired-state!) (let ([result '()]) (persistent-map-for-each (lambda (id record) (set! result (cons record result))) @@ -685,9 +708,15 @@ (cond [(or (hash-key? *replay-keys* event-key) (hash-key? *replay-keys* nonce-key)) #f] - [(> (+ (hashtable-size *replay-keys*) 2) *max-replay-keys*) #f] + ;; Distinguish capacity exhaustion from a genuine replay so the caller + ;; can answer 503 quota instead of a misleading 409 "replayed event". + [(> (+ (hashtable-size *replay-keys*) 2) *max-replay-keys*) 'quota] [(not (db-claim-replay! event-key expires now)) #f] - [(not (db-claim-replay! nonce-key expires now)) #f] + [(not (db-claim-replay! nonce-key expires now)) + ;; Roll back the event-key claim so a partial two-key claim never + ;; leaves the DB key claimed while the in-memory cache is absent. + (db-release-replay! event-key) + #f] [else (hash-put! *replay-keys* event-key expires) (hash-put! *replay-keys* nonce-key expires) @@ -1281,24 +1310,29 @@ [(not (and (safe-token? (hash-get payload "id") 1 256) (string=? (hash-get payload "id") event-id))) (respond-json 400 "{\"error\":\"signed event id must match body id\"}")] - [(not (claim-replay-key! source event-id nonce timestamp)) - (respond-json 409 "{\"error\":\"replayed event\"}")] - [else - (let ([event (list->hash-table - `(("id" . ,event-id) - ("type" . ,type) - ("source" . ,source) - ("payload" . ,payload) - ("received" . ,(format "~a" (time-second (current-time)))) - ("status" . "queued")))]) - (if (chan-try-put! ingest-ch event) - (begin - (counter-inc! *metric-received*) - (respond-json 202 - (json-object->string - (list->hash-table - `(("status" . "accepted") ("id" . ,event-id)))))) - (respond-json 503 "{\"error\":\"ingest queue full\"}")))]))]))) + [else + (let ([claim (claim-replay-key! source event-id nonce timestamp)]) + (cond + [(eq? claim 'quota) + (respond-json 503 "{\"error\":\"quota exceeded\"}")] + [(not claim) + (respond-json 409 "{\"error\":\"replayed event\"}")] + [else + (let ([event (list->hash-table + `(("id" . ,event-id) + ("type" . ,type) + ("source" . ,source) + ("payload" . ,payload) + ("received" . ,(format "~a" (time-second (current-time)))) + ("status" . "queued")))]) + (if (chan-try-put! ingest-ch event) + (begin + (counter-inc! *metric-received*) + (respond-json 202 + (json-object->string + (list->hash-table + `(("status" . "accepted") ("id" . ,event-id)))))) + (respond-json 503 "{\"error\":\"ingest queue full\"}")))]))]))]))) ;; GET /api/events/:id — query a processed event (def (handle-get-event req) @@ -1488,15 +1522,17 @@ (respond-json 404 "{\"error\":\"not found in dead-letter store\"}")] [else ;; Re-enqueue with reset retry count - (let ([replay-event - (list->hash-table - `(("id" . ,id) - ("type" . ,(hash-ref record "type" "unknown")) - ("source" . "edge-operator") - ("payload" . ,(make-hash-table)) - ("received" . ,(format "~a" (time-second (current-time)))) - ("retry-count" . 0) - ("status" . "replayed")))]) + (let ([replay-event + (list->hash-table + `(("id" . ,id) + ("type" . ,(hash-ref record "type" "unknown")) + ("source" . "edge-operator") + ("payload" . ,(or (parse-json-object + (hash-ref record "payload" "{}")) + (make-hash-table))) + ("received" . ,(format "~a" (time-second (current-time)))) + ("retry-count" . 0) + ("status" . "replayed")))]) (let ([queued? (with-mutex *state-mutex* (let ([st (ref-deref dead-letter-ref)]) @@ -1538,6 +1574,7 @@ (guard (e [(condition? e) (void)]) (fiber-httpd-stop! *httpd-server*))) ;; Convert delayed work to bounded dead-letter state before closing ingest. (stop-retry-scheduler!) + (stop-prune-scheduler!) ;; Close ingest channel — workers' (<!! ingest-ch) returns EOF (guard (e [(condition? e) (void)]) (chan-close! ingest-ch)) ;; Drain window — let in-flight events finish @@ -1581,6 +1618,7 @@ (db-init!) (db-load-events!) (start-retry-scheduler!) + (start-prune-scheduler!) ;; Register SIGTERM handler (Phase 3.6) (add-signal-handler! SIGTERM graceful-shutdown!)