query/group: native columnar group-by / aggregate

ober

f70f0ebc38d38b913505d6d0960b33d0eebde4db

diff --git a/lib/jerboa-db/query/group.ss b/lib/jerboa-db/query/group.ss
new file mode 100644
index 0000000..78d70fb
--- /dev/null
+++ b/lib/jerboa-db/query/group.ss
@@ -0,0 +1,122 @@
+#!chezscheme
+;;; (jerboa-db query group) — Native columnar group-by / aggregate
+;;;
+;;; A direct, vectorised group-by-aggregate that bypasses the Datalog binding
+;;; machinery for the common analytic shape "aggregate VALUE-ATTR grouped by
+;;; GROUP-ATTR" (both cardinality-one attributes of the same entity). This is
+;;; the in-process answer to the queries that currently get routed to DuckDB
+;;; (Q5 "releases per country", Q8 "avg track duration by status").
+;;;
+;;; How it works:
+;;;   1. Scan AEVT for each attribute (one contiguous range), resolve to the
+;;;      current value per entity (same rule as the engine).
+;;;   2. Build the value attribute's current datoms into a columnar segment, so
+;;;      the aggregate loop reads an unboxed primitive value column.
+;;;   3. Join group/value by entity via an e->group hashtable and hash-aggregate.
+;;;
+;;; Scope: single-hop (group and value on the same entity), cardinality-one.
+;;; Multi-hop grouping (e.g. via a ref) and auto-routing from `q` are follow-ons.
+
+(library (jerboa-db query group)
+  (export group-aggregate group-count)
+
+  (import (except (chezscheme)
+                  make-hash-table hash-table?
+                  sort sort!
+                  printf fprintf
+                  path-extension path-absolute?
+                  with-input-from-string with-output-to-string
+                  iota 1+ 1-
+                  partition
+                  make-date make-time
+                atom? meta)
+          (jerboa prelude)
+          (jerboa-db datom)
+          (jerboa-db schema)
+          (jerboa-db history)
+          (jerboa-db index protocol)
+          (jerboa-db segment))
+
+  ;; Resolve a scanned AEVT range to current datoms (mirrors the engine:
+  ;; all-asserted -> pass through; otherwise keep the latest (e,a,v) and drop
+  ;; those whose latest is a retraction).
+  (def (resolve-current datoms)
+    (if (for-all datom-added? datoms)
+        datoms
+        (let ([ht (make-hashtable equal-hash equal?)])
+          (for-each
+            (lambda (d)
+              (let ([key (list (datom-e d) (datom-a d) (datom-v d))])
+                (let ([ex (hashtable-ref ht key #f)])
+                  (when (or (not ex) (> (datom-tx d) (datom-tx ex)))
+                    (hashtable-set! ht key d)))))
+            datoms)
+          (let-values ([(ks vs) (hashtable-entries ht)])
+            (filter datom-added? (vector->list vs))))))
+
+  ;; Current datoms for one attribute id, honouring db-level (as-of/since) filters.
+  (def (current-datoms-for db aid)
+    (let* ([aevt (db-resolve-index db 'aevt)]
+           [lo (make-datom 0 aid +min-val+ 0 #t)]
+           [hi (make-datom (greatest-fixnum) aid +max-val+ (greatest-fixnum) #t)]
+           [raw (dbi-range aevt lo hi)]
+           [filtered (filter (lambda (d) (db-filter-datom? db d)) raw)])
+      (resolve-current filtered)))
+
+  (def (attr-id db ident)
+    (let ([attr (schema-lookup-by-ident (db-value-schema db) ident)])
+      (if attr (db-attribute-id attr)
+          (error 'group-aggregate "unknown attribute" ident))))
+
+  ;; cell = #(count sum min max); v may be #f (count-only)
+  (def (cell-bump! acc g v)
+    (let ([cell (hashtable-ref acc g #f)])
+      (if cell
+          (begin
+            (vector-set! cell 0 (+ (vector-ref cell 0) 1))
+            (when v
+              (vector-set! cell 1 (+ (vector-ref cell 1) v))
+              (when (< v (vector-ref cell 2)) (vector-set! cell 2 v))
+              (when (> v (vector-ref cell 3)) (vector-set! cell 3 v))))
+          (hashtable-set! acc g (vector 1 (or v 0) (or v 0) (or v 0))))))
+
+  (def (project-agg agg cell)
+    (let ([cnt (vector-ref cell 0)] [sum (vector-ref cell 1)]
+          [mn (vector-ref cell 2)] [mx (vector-ref cell 3)])
+      (case agg
+        [(count) cnt]
+        [(sum)   sum]
+        [(avg)   (if (= cnt 0) 0 (/ sum cnt))]
+        [(min)   mn]
+        [(max)   mx]
+        [else (error 'group-aggregate "unknown aggregate" agg)])))
+
+  ;; Aggregate value-attr grouped by group-attr. value-attr may be #f for a
+  ;; plain count of entities per group (agg must then be 'count). Returns an
+  ;; alist (group-value . result).
+  (def (group-aggregate db group-attr value-attr agg)
+    (let ([acc (make-hashtable equal-hash equal?)])  ;; group value -> cell
+      (if value-attr
+          ;; join group/value by entity, then aggregate the value column unboxed
+          (let ([ght (make-hashtable equal-hash equal?)])  ;; entity -> group value
+            (for-each (lambda (d) (hashtable-set! ght (datom-e d) (datom-v d)))
+                      (current-datoms-for db (attr-id db group-attr)))
+            (let* ([vseg (make-segment (current-datoms-for db (attr-id db value-attr)))]
+                   [n (segment-count vseg)])
+              (let loop ([i 0])
+                (when (< i n)
+                  (let ([g (hashtable-ref ght (segment-e vseg i) #f)])
+                    (when g (cell-bump! acc g (segment-v vseg i))))
+                  (loop (+ i 1))))))
+          ;; count-only: tally directly by group value (no entity table)
+          (for-each (lambda (d) (cell-bump! acc (datom-v d) #f))
+                    (current-datoms-for db (attr-id db group-attr))))
+      (let-values ([(ks vs) (hashtable-entries acc)])
+        (map (lambda (g cell) (cons g (project-agg agg cell)))
+             (vector->list ks) (vector->list vs)))))
+
+  ;; Count of entities per distinct value of group-attr.
+  (def (group-count db group-attr)
+    (group-aggregate db group-attr #f 'count))
+
+) ;; end library
diff --git a/tests/test-core.ss b/tests/test-core.ss
index 7f6ac92..6da1d0c 100644
--- a/tests/test-core.ss
+++ b/tests/test-core.ss
@@ -8,7 +8,8 @@
         (jerboa-db value-store)
         (jerboa-db index protocol)
         (jerboa-db index memory)
-        (jerboa-db segment))
+        (jerboa-db segment)
+        (jerboa-db query group))
 
 ;; ---- Test harness ----
 
@@ -736,6 +737,22 @@
     (assert-equal (segment-added? seg2 1) #f)
     (assert-equal (segment-v seg2 1) "beta")))
 
+(test "native columnar group-by aggregates correctly"
+  (let ([conn (connect ":memory:")]
+        [sort-alist (lambda (al) (list-sort (lambda (x y) (string<? (car x) (car y))) al))])
+    (transact! conn
+      (list '((db/ident . g/cat) (db/valueType . db.type/string) (db/cardinality . db.cardinality/one))
+            '((db/ident . g/val) (db/valueType . db.type/long) (db/cardinality . db.cardinality/one))))
+    (transact! conn
+      (list '((g/cat . "a") (g/val . 10))
+            '((g/cat . "a") (g/val . 20))
+            '((g/cat . "b") (g/val . 100))))
+    (let ([d (db conn)])
+      (assert-equal (sort-alist (group-count d 'g/cat)) '(("a" . 2) ("b" . 1)))
+      (assert-equal (sort-alist (group-aggregate d 'g/cat 'g/val 'sum)) '(("a" . 30) ("b" . 100)))
+      (assert-equal (sort-alist (group-aggregate d 'g/cat 'g/val 'avg)) '(("a" . 15) ("b" . 100)))
+      (assert-equal (sort-alist (group-aggregate d 'g/cat 'g/val 'max)) '(("a" . 20) ("b" . 100))))))
+
 ;; ============================================================
 ;; Report
 ;; ============================================================