query: add :limit with streaming early-exit
ober
2c904a0d6a379b6a45bfac79c0840a66edc96b92
--- a/lib/jerboa-db/query/engine.ss +++ b/lib/jerboa-db/query/engine.ss @@ -111,26 +111,30 @@ (fields find-vars ;; list of symbols or (agg ?var) forms in-vars ;; list of input symbols ($ = db) where-clauses ;; list of clause forms - rules)) + rules + limit)) ;; integer row cap, or #f (def (parse-query form) (let ([find-clause #f] [in-clause '($)] [where-clause '()] - [rule-defs #f]) + [rule-defs #f] + [limit-val #f]) (for-each (lambda (section) (case (car section) [(find) (set! find-clause (cdr section))] [(in) (set! in-clause (cdr section))] [(where) (set! where-clause (cdr section))] + [(limit) (when (pair? (cdr section)) (set! limit-val (cadr section)))] [else (void)])) form) (make-parsed-query (or find-clause '()) in-clause where-clause - rule-defs))) + rule-defs + limit-val))) ;; ---- Current-state resolution ---- ;; In a non-history view, datoms are append-only: both assertions (added?=#t) @@ -1162,9 +1166,111 @@ find-vars))) alist))) - ;; ---- Top-level query function ---- + ;; ---- Streaming :limit (early-exit) ---- - (def (query-db parsed db . inputs) + (def (limit-take lst n) + (cond + [(not n) lst] + [(<= n 0) '()] + [else (let loop ([l lst] [k n] [acc '()]) + (if (or (null? l) (= k 0)) + (reverse acc) + (loop (cdr l) (- k 1) (cons (car l) acc))))])) + + (def (data-pattern-clause? c) + (and (pair? c) (not (pair? (car c))) (>= (length c) 3) + (symbol? (cadr c)) (not (logic-var? (cadr c))))) + + ;; Streaming early-exit is sound only for current (non-history) views, no input + ;; bindings, and non-aggregate finds (aggregates need every row). + (def (limit-streamable? db find-vars in-vars) + (and (not (db-value-history? db)) + (equal? in-vars '($)) + (not (any-aggregates? find-vars)))) + + ;; Stream the current solutions of a single data pattern from empty bindings, + ;; calling (yield bindings) for each. Uses a cursor (lazy) + run-grouping to + ;; resolve current state, so an early escape stops the index scan. Non-history. + (def (scan-pattern-current db pattern schema yield) + (let* ([e-spec (car pattern)] + [a-spec (cadr pattern)] + [v-spec (if (>= (length pattern) 3) (caddr pattern) '_)] + [tx-spec (if (>= (length pattern) 4) (cadddr pattern) '_)] + [op-spec (if (>= (length pattern) 5) (list-ref pattern 4) '_)] + [empty (make-empty-bindings)] + [e-val (resolve-in-bindings empty e-spec)] + [v-val (resolve-in-bindings empty v-spec)] + [attr (schema-lookup-by-ident schema a-spec)]) + (when attr + (let* ([aid (db-attribute-id attr)] + [index-name (cond + [e-val 'eavt] + [(and v-val (ref-type? attr)) 'vaet] + [(and v-val (avet-eligible? attr)) 'avet] + [else 'aevt])] + [idx (db-resolve-index db index-name)] + [lo (case index-name + [(eavt aevt) (make-datom (or e-val 0) aid (if v-val v-val +min-val+) 0 #t)] + [(avet vaet) (make-datom 0 aid (if v-val v-val +min-val+) 0 #t)])] + [hi (case index-name + [(eavt aevt) (make-datom (or e-val (greatest-fixnum)) aid + (if v-val v-val +max-val+) (greatest-fixnum) #t)] + [(avet vaet) (make-datom (greatest-fixnum) aid + (if v-val v-val +max-val+) (greatest-fixnum) #t)])]) + (define (same-eav? a b) + (and (= (datom-e a) (datom-e b)) (= (datom-a a) (datom-a b)) + (equal? (datom-v a) (datom-v b)))) + (define (emit run-last) + (when (and run-last (datom-added? run-last)) + (let ([b (try-unify-datom run-last e-spec a-spec v-spec tx-spec op-spec empty)]) + (when b (yield b))))) + ;; Walk the cursor; each maximal (e,a,v) run's last db-filtered datom is + ;; that triple's current state (tx is the final sort key, ascending). + (let loop ([s (dbi-cursor idx lo hi)] [run-last #f]) + (if (stream-null? s) + (emit run-last) + (let ([d (stream-car s)]) + (if (db-filter-datom? db d) + (if (and run-last (same-eav? run-last d)) + (loop (stream-cdr s) d) + (begin (emit run-last) (loop (stream-cdr s) d))) + (loop (stream-cdr s) run-last))))))))) + + ;; Evaluate a non-aggregate query with a row limit, stopping the scan as soon + ;; as `limit` distinct rows are found: stream the (most selective) first clause + ;; and evaluate the rest eagerly per binding. Falls back to the caller for + ;; shapes this can't stream. + (def (eval-limit db ordered-clauses find-vars schema rules-ht limit) + (let ([first-clause (car ordered-clauses)] + [rest-clauses (cdr ordered-clauses)] + [seen (make-hashtable equal-hash equal?)] + [acc '()] + [n 0]) + (let ([rest-live (compute-live-vars-at-each-step rest-clauses find-vars)]) + (call/cc + (lambda (stop) + (when (<= limit 0) (stop #f)) + (scan-pattern-current db first-clause schema + (lambda (binding) + (let ([tails (if (null? rest-clauses) + (list binding) + (evaluate-where-clauses db rest-clauses (list binding) + schema rules-ht rest-live))]) + (for-each + (lambda (b) + (let ([row (map (lambda (v) (if (logic-var? v) (binding-ref b v) v)) + find-vars)]) + (unless (hashtable-ref seen row #f) + (hashtable-set! seen row #t) + (set! acc (cons row acc)) + (set! n (+ n 1)) + (when (>= n limit) (stop #f))))) + tails)))))) + (reverse acc)))) + + ;; Full (unlimited) routing: pure-aggregate fast path, native group-by, DuckDB + ;; analytics fallback, else the general Datalog engine. + (def (routed-result parsed db inputs) (let* ([schema (db-value-schema db)] [find-vars (parsed-query-find-vars parsed)] [in-vars (parsed-query-in-vars parsed)] @@ -1174,8 +1280,6 @@ (cond [pure-agg-desc (execute-pure-aggregate db pure-agg-desc)] - ;; Native columnar group-by (single-hop / one ref-hop) — beats the DuckDB - ;; round-trip and the general join, with identical results. [(and (not (db-value-history? db)) (native-group-descriptor schema find-vars in-vars where-clauses)) => (lambda (desc) (execute-native-group db desc find-vars))] @@ -1190,6 +1294,30 @@ (query-db-general schema find-vars in-vars where-clauses rules-ht db inputs)]))) + ;; ---- Top-level query function ---- + + (def (query-db parsed db . inputs) + (let ([find-vars (parsed-query-find-vars parsed)] + [in-vars (parsed-query-in-vars parsed)] + [where-clauses (parsed-query-where-clauses parsed)] + [limit (parsed-query-limit parsed)]) + (cond + ;; :limit with early-exit — stream the most-selective first data pattern, + ;; evaluate the rest eagerly per binding, stop at `limit` distinct rows. + [(and limit (pair? where-clauses) (limit-streamable? db find-vars in-vars)) + (let* ([schema (db-value-schema db)] + [db-stats (db-value-stats db)] + [card-fn (make-card-fn db)] + [ordered (reorder-clauses where-clauses '() schema db-stats card-fn)]) + (if (and (pair? ordered) (data-pattern-clause? (car ordered))) + (eval-limit db ordered find-vars schema + (parsed-query-rules parsed) limit) + (limit-take (routed-result parsed db inputs) limit)))] + ;; :limit on a shape we don't stream — correct via truncation + [limit (limit-take (routed-result parsed db inputs) limit)] + ;; no limit — full routing + [else (routed-result parsed db inputs)]))) + (def (query-db-general schema find-vars in-vars where-clauses rules-ht db inputs) (let* ([rules-ht rules-ht] ;; Build initial bindings from inputs --- a/tests/test-core.ss +++ b/tests/test-core.ss @@ -857,6 +857,33 @@ (assert-equal (length rows) 1) (assert-equal (caar rows) eid1))))) +(test "query :limit caps results (streaming + truncation paths)" + (let ([conn (connect ":memory:")]) + (transact! conn + (list '((db/ident . t/kind) (db/valueType . db.type/string) (db/cardinality . db.cardinality/one)) + '((db/ident . t/n) (db/valueType . db.type/long) (db/cardinality . db.cardinality/one)))) + (transact! conn + (for/collect ([i (in-range 20)]) + `((t/kind . ,(if (even? i) "a" "b")) (t/n . ,i)))) + (let ([d (db conn)]) + ;; single-pattern limit (streaming early-exit): 20 match, take 5 + (assert-equal (length (q '((find ?e) (where (?e t/n ?x)) (limit 5)) d)) 5) + ;; limit larger than total -> all rows + (assert-equal (length (q '((find ?e) (where (?e t/kind "a")) (limit 100)) d)) 10) + ;; limit 0 -> empty + (assert-equal (q '((find ?e) (where (?e t/n ?x)) (limit 0)) d) '()) + ;; multi-clause limit (stream first clause, eager rest): 10 candidates, take 3 + (assert-equal (length (q '((find ?e ?x) (where (?e t/kind "a") (?e t/n ?x)) (limit 3)) d)) 3) + ;; aggregate + limit (truncation path): 2 groups, capped to 1 + (assert-equal (length (q '((find ?k (count ?e)) (where (?e t/kind ?k) (?e t/n ?x)) (limit 1)) d)) 1) + ;; same aggregate without limit -> 2 groups + (assert-equal (length (q '((find ?k (count ?e)) (where (?e t/kind ?k) (?e t/n ?x))) d)) 2) + ;; every limited row is a genuine solution (?e really has kind "a") + (let* ([all-a (map car (q '((find ?e) (where (?e t/kind "a"))) d))] + [lim (map car (q '((find ?e) (where (?e t/kind "a")) (limit 4)) d))]) + (assert-equal (length lim) 4) + (assert-true (for-all (lambda (e) (memv e all-a)) lim)))))) + ;; ============================================================ ;; Report ;; ============================================================