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
6 changes: 4 additions & 2 deletions tools/codegen_jvm.py
Original file line number Diff line number Diff line change
Expand Up @@ -690,13 +690,15 @@ def _codecs(self):
self.value_class[sql] = cls

def _enum_parsers(self):
"""For each C enum, the catalog function reading it from its name."""
"""For each C enum, the public catalog function reading it from its name, as a binding
calls the public API alone; the Spark arm reads the same (ENUM_PARSER in
codegen_spark_udfs.py)."""
self.enum_parser = {}
for f in self.fns:
rt = _norm(f['returnType']['canonical'])
ps = f['params']
if rt in self.enums and len(ps) == 1 and _norm(ps[0]['canonical']) == 'char *' \
and f['name'] in self.jmeos:
and f.get('api') == 'public' and f['name'] in self.jmeos:
self.enum_parser.setdefault(rt, f['name'])


Expand Down
32 changes: 25 additions & 7 deletions tools/codegen_spark_udfs.py
Original file line number Diff line number Diff line change
Expand Up @@ -202,6 +202,9 @@ def arg_kind(canon):
return ("ts",)
if b in PARSE:
return ("ptr",) + PARSE[b]
# An enum the catalog reads from its name: the SQL text is parsed into the int JMEOS takes.
if b in ENUM_PARSER and "*" not in nc:
return ("scalar", "StringType", "String", "GeneratedFunctions.%s(%%s)" % ENUM_PARSER[b])
# scalar ONLY when not a pointer: int* / DateADT* are arrays/out-params, not ints.
if b in SCALAR_ARG and "*" not in nc:
return ("scalar",) + SCALAR_ARG[b]
Expand Down Expand Up @@ -267,8 +270,10 @@ def ret_emit(canon, sqlop):

def supported(f):
"""Reason string if NOT emittable, else None."""
# meos_internal_* doxygen groups are MEOS-internal, not user-facing — excluded.
if (f.get("group") or "").startswith("meos_internal"):
# A binding calls the public API alone. The catalog's `api` states it, public for a
# function whose @ingroup is a public group; a function stating no group reads as
# internal there, so its name is no test of it.
if f.get("api") != "public":
return "internal"
in_params, out = classify(f)
if out is None:
Expand Down Expand Up @@ -452,7 +457,7 @@ def emit_single(name, f, vis_arity=None):
L.append(f" java.time.OffsetDateTime dt_{a} = UdfMarshal.tsOdt({a});")
callargs.append(f"dt_{a}")
else:
callargs.append(a)
callargs.append(k[3] % a)
# supply the wrapper-bound literal (shape.boundArgs) — or the generic type default —
# for the SQL-hidden trailing flags (sqlArity..C-arity)
for p in hidden:
Expand Down Expand Up @@ -607,6 +612,10 @@ def _famrank(f):
# that concrete temporal type, so the dispatcher can tell it apart from a sibling overload
# by the WKB type byte instead of guessing.
TEMPTYPE_CODE = {}
# For each C enum, the public catalog function reading it from its name, so an enum argument
# travels as the text its SQL function takes. The rule is the Flink arm's, SqlModel._enum_parsers
# in codegen_jvm.py, read from the public functions alone, as #supported admits only those.
ENUM_PARSER = {}
# The value of each macro and enum member the catalog states, which is what a bound
# literal or SQL default naming one passes. Filled from the catalog before any emit pass.
CONST = {}
Expand Down Expand Up @@ -699,9 +708,9 @@ def permissiveness(f):
callargs.append("D_%s" % a)
elif i in mixed:
classes.append("%s instanceof %s" % (a, k[2]))
callargs.append("((%s) %s)" % (k[2], a))
callargs.append(k[3] % ("((%s) %s)" % (k[2], a)))
else:
callargs.append(a)
callargs.append(k[3] % a)
# SQL-hidden trailing flags get the wrapper-bound literal (shape.boundArgs) or the
# generic default (zip above paired only the first `vis` exposed args; the
# candidate's remaining params are the flags).
Expand Down Expand Up @@ -903,7 +912,7 @@ def emit_scalar_values(name, f, shape):
L.append(" java.time.OffsetDateTime dt_%s = UdfMarshal.tsOdt(%s);" % (a, a))
callargs.append("dt_%s" % a)
else:
callargs.append(a)
callargs.append(k[3] % a)
L.append(" jnr.ffi.Runtime _rt = jnr.ffi.Runtime.getSystemRuntime();")
L.append(" jnr.ffi.Pointer _cnt = jnr.ffi.Memory.allocateDirect(_rt, 4);")
callargs.append("_cnt")
Expand Down Expand Up @@ -1116,7 +1125,7 @@ def emit_setret(name, cands, colnames):
elif k[0] == "ts":
callargs.append("UdfMarshal.tsOdt(%s)" % a)
else:
callargs.append(a)
callargs.append(k[3] % a)
L.append(" if (%s) {" % (" && ".join("%s != null" % p for p in ptrs) or "true"))
L.append(" jnr.ffi.Runtime _rt = jnr.ffi.Runtime.getSystemRuntime();")
cells = {}
Expand Down Expand Up @@ -1661,6 +1670,15 @@ def main():
0 if not jargs else len(jargs.split(",")))
# The codec each value travels in, from the catalog, before any emit pass reads it.
derive_codecs(cat, lambda n: bool(n) and (jar_syms is None or n in jar_syms))
# The public parser of each enum the catalog states, before any emit pass reads it.
enums = {e["name"] for e in cat.get("enums", [])}
for f in fns:
ps = f["params"]
rt = norm(f["returnType"]["canonical"])
if (rt in enums and len(ps) == 1 and norm(ps[0]["canonical"]) == "char *"
and f.get("api") == "public"
and (jar_syms is None or f["name"] in jar_syms)):
ENUM_PARSER.setdefault(rt, f["name"])

# GOAL: reach the WHOLE JMEOS surface. Every MEOS C function (unique by its C
# name) becomes a 1:1 UDF named by that C symbol — that is how the ~2254
Expand Down
33 changes: 0 additions & 33 deletions tools/spark-udf-gaps.txt
Original file line number Diff line number Diff line change
Expand Up @@ -138,7 +138,6 @@ index_result_create no-encoder:MeosArray
index_result_id no-decoder:MeosArray; array-or-out-param:id
int32_hash ret:uint32_t
int64_hash ret:uint32_t
interptype_from_string internal
intersection_posechain_set arg:PoseChain *
intersection_set_posechain arg:PoseChain *
intersection_set_text internal
Expand Down Expand Up @@ -193,7 +192,6 @@ jsonb_set array-or-out-param:path_elems
jsonb_set_lax array-or-out-param:path_elems
jsonb_to_float4 ret:float
jsonb_to_int16 ret:int16_t
jsonbset_array_element arg:nullHandleType
jsonbset_delete internal
jsonbset_delete_array array-or-out-param:keys
jsonbset_delete_path array-or-out-param:path_elems
Expand All @@ -207,10 +205,6 @@ jsonbset_path_exists array-or-out-param:count; unsupported-return:bool *
jsonbset_path_match array-or-out-param:count; unsupported-return:bool *
jsonbset_set array-or-out-param:keys
jsonbset_to_alphanumset internal
jsonbset_to_bigintset arg:nullHandleType
jsonbset_to_floatset arg:nullHandleType
jsonbset_to_intset arg:nullHandleType
jsonbset_to_textset_key arg:nullHandleType
jsonbset_value_n internal
le_date_timestamp arg:Timestamp
le_timestamp_date arg:Timestamp
Expand Down Expand Up @@ -243,7 +237,6 @@ ne_timestamptz_timestamp arg:Timestamp
npoint_hash ret:uint32_t
npointset_make internal
npointset_value_n internal
null_handle_type_from_string ret:nullHandleType
overabove_tpcbox_tpcbox arg:TPCBox *
overafter_tpcbox_tpcbox arg:TPCBox *
overback_tpcbox_tpcbox arg:TPCBox *
Expand Down Expand Up @@ -446,25 +439,18 @@ tbox_tmin arg:TimestampTz *
tcbuffer_value_at_timestamptz internal
tcbuffer_value_n internal
tcbuffer_values internal
tcbufferseq_from_base_tstzspan internal
tcbufferseqset_from_base_tstzspanset internal
temparr_round unsupported-return:Temporal **
temporal_append_tinstant internal
temporal_as_tsequence internal
temporal_as_tsequenceset internal
temporal_hash ret:uint32_t
temporal_instants internal
temporal_merge_array internal
temporal_segments internal
temporal_sequences internal
temporal_set_interp internal
temporal_spans array-or-out-param:count
temporal_split_each_n_spans array-or-out-param:count
temporal_split_n_spans array-or-out-param:count
temporal_time_bins array-or-out-param:count
temporal_timestamps array-or-out-param:count
temporal_timestamptz_n arg:TimestampTz *
temporal_tsample internal
teq_posechain_tposechain arg:PoseChain *
teq_text_ttext internal
teq_tposechain_posechain arg:PoseChain *
Expand Down Expand Up @@ -501,8 +487,6 @@ tfloat_value_time_boxes array-or-out-param:count
tfloat_wmax_transfn no-decoder:SkipList; no-encoder:SkipList
tfloat_wmin_transfn no-decoder:SkipList; no-encoder:SkipList
tfloat_wsum_transfn no-decoder:SkipList; no-encoder:SkipList
tfloatseq_from_base_tstzspan internal
tfloatseqset_from_base_tstzspanset internal
tge_text_ttext internal
tge_ttext_text internal
tgeo_space_boxes array-or-out-param:count
Expand All @@ -519,8 +503,6 @@ tgeogpoint_s2cell_split array-or-out-param:cells
tgeompoint_h3index_split array-or-out-param:cells
tgeompoint_quadbin_split array-or-out-param:cells
tgeompoint_s2cell_split array-or-out-param:cells
tgeoseq_from_base_tstzspan internal
tgeoseqset_from_base_tstzspanset internal
tgt_text_ttext internal
tgt_ttext_text internal
th3index_value_at_timestamptz arg:uint64_t *
Expand Down Expand Up @@ -589,10 +571,8 @@ tint_value_time_boxes array-or-out-param:count
tint_wmax_transfn no-decoder:SkipList; no-encoder:SkipList
tint_wmin_transfn no-decoder:SkipList; no-encoder:SkipList
tint_wsum_transfn no-decoder:SkipList; no-encoder:SkipList
tjson_array_element arg:nullHandleType
tjson_extract_path array-or-out-param:path_elems
tjson_object_field internal
tjsonb_array_element arg:nullHandleType
tjsonb_delete internal
tjsonb_delete_array array-or-out-param:keys
tjsonb_delete_path array-or-out-param:path_elems
Expand All @@ -604,11 +584,6 @@ tjsonb_extract_path array-or-out-param:path_elems
tjsonb_insert array-or-out-param:keys
tjsonb_object_field internal
tjsonb_set array-or-out-param:keys
tjsonb_to_tbigint arg:nullHandleType
tjsonb_to_tbool arg:nullHandleType
tjsonb_to_tfloat internal
tjsonb_to_tint arg:nullHandleType
tjsonb_to_ttext_key arg:nullHandleType
tjsonb_value_at_timestamptz internal
tjsonb_value_n internal
tjsonb_values internal
Expand All @@ -625,8 +600,6 @@ tnpoint_tcentroid_transfn no-decoder:SkipList; no-encoder:SkipList
tnpoint_value_at_timestamptz internal
tnpoint_value_n internal
tnpoint_values internal
tnpointseq_from_base_tstzspan internal
tnpointseqset_from_base_tstzspanset internal
tnumber_split_each_n_tboxes array-or-out-param:count
tnumber_split_n_tboxes array-or-out-param:count
tnumber_tavg_combinefn no-decoder:SkipList; no-encoder:SkipList
Expand Down Expand Up @@ -679,9 +652,7 @@ tpoint_make_simple internal
tpoint_tcentroid_finalfn no-decoder:SkipList
tpoint_tcentroid_transfn no-decoder:SkipList; no-encoder:SkipList
tpoint_tfloat_to_geomeas internal
tpointseq_from_base_tstzspan internal
tpointseq_make_coords array-or-out-param:xcoords; array-or-out-param:ycoords; array-or-out-param:zcoords
tpointseqset_from_base_tstzspanset internal
tpose_value_at_timestamptz internal
tpose_value_n internal
tpose_values internal
Expand All @@ -693,18 +664,14 @@ tposechaininst_make arg:PoseChain *
tposechainseq_from_base_tstzset arg:PoseChain *
tposechainseq_from_base_tstzspan arg:PoseChain *
tposechainseqset_from_base_tstzspanset arg:PoseChain *
tposeseq_from_base_tstzspan internal
tposeseqset_from_base_tstzspanset internal
tquadbin_value_at_timestamptz arg:uint64_t *
tquadbin_value_n arg:uint64_t *
tquadbinseq_make array-or-out-param:values
tquadbinseqset_make array-or-out-param:sequences
trgeometry_append_tinstant internal
trgeometry_instants internal
trgeometry_merge_array internal
trgeometry_segments internal
trgeometry_sequences internal
trgeometry_set_interp internal
trgeometry_space_boxes array-or-out-param:count
trgeometry_space_time_boxes array-or-out-param:count
trgeometry_split_each_n_stboxes array-or-out-param:count
Expand Down
Loading