This is an automated email from the ASF dual-hosted git repository.

Tartarus0zm pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/auron.git


The following commit(s) were added to refs/heads/master by this push:
     new 83f05f58a [AURON #1863] Support native Flink UNIX_TIMESTAMP: native 
scalar function (#2409)
83f05f58a is described below

commit 83f05f58ac0bc5d9497c5946e9386758b0118b94
Author: Weiqing Yang <[email protected]>
AuthorDate: Tue Aug 4 02:49:40 2026 -0700

    [AURON #1863] Support native Flink UNIX_TIMESTAMP: native scalar function 
(#2409)
    
    # Which issue does this PR close?
    
    Closes #1863. The converter that emits this function follows in a
    separate PR, once this one merges.
    
    To show where this is going, the converter change is previewable here:
    https://github.com/weiqingy/auron/pull/1 (a draft in my fork, diffed
    against this PR's branch so it shows only the converter side). That PR
    will be opened against this repo after this one lands.
    
    # Rationale for this change
    
    Flink's `UNIX_TIMESTAMP` converts a formatted date-time string to a Unix
    timestamp in seconds. AIP-1 lists it as a Phase 1 built-in function.
    Supporting it natively is the first step; this PR adds the native (Rust)
    evaluation, and a follow-up PR wires the Flink converter to emit it.
    
    The hard part is matching Flink's parsing exactly. Flink parses with
    `java.text.SimpleDateFormat` in its default lenient mode, so it accepts
    non-zero-padded fields, field rollover, and trailing input. A strict
    parser diverges from Flink on the default format `yyyy-MM-dd HH:mm:ss`
    as soon as a field is not zero-padded, and would return a wrong
    timestamp rather than an error. So this implements a lenient parser
    rather than delegating to a strict one.
    
    # What changes are included in this PR?
    
    A new `Flink_UnixTimestamp` ext scalar function in
    `datafusion-ext-functions`, registered through the existing ext-function
    path (no proto change).
    
    It takes three string arguments `[value, format, zoneId]` and returns an
    `Int64`. The format is a translated strftime-style pattern and the zone
    is resolved by the caller; the converter PR supplies both.
    
    The parser reproduces `SimpleDateFormat` lenient semantics for the
    supported fields: variable field widths, rollover normalization,
    tolerant of trailing input, a leading minus on numeric fields. It also
    handles the Julian/Gregorian hybrid calendar (cutover 1582-10-15), since
    `GregorianCalendar` is a hybrid while chrono is proleptic Gregorian, and
    resolves DST-ambiguous or nonexistent local times using the zone's
    standard offset, matching Flink.
    
    Failure semantics match Flink: an unparseable string yields
    `Long.MIN_VALUE`, a NULL input yields NULL, and the
    milliseconds-to-seconds step truncates toward zero.
    
    The function is registered but not yet emitted by any planner, so this
    PR is self-contained and dormant until the converter lands.
    
    # Are there any user-facing changes?
    
    No. The function is not reachable from SQL until the converter PR is
    merged.
    
    # How was this patch tested?
    
    19 unit tests derived from a differential oracle that was validated
    against real `java.text.SimpleDateFormat` across a large randomized set
    of (input, format, timezone) inputs. Coverage includes the happy path,
    field rollover, non-padded fields, trailing garbage, NULL,
    unparseable-to-`Long.MIN_VALUE`, DST gap and overlap across several
    zones, pre-1970 negative timestamps, pre-1582 hybrid-calendar dates, the
    arity guard, an empty format matching nothing, and canonical field
    widths.
---
 .../datafusion-ext-functions/src/flink_datetime.rs | 668 +++++++++++++++++++++
 native-engine/datafusion-ext-functions/src/lib.rs  |   4 +-
 2 files changed, 671 insertions(+), 1 deletion(-)

diff --git a/native-engine/datafusion-ext-functions/src/flink_datetime.rs 
b/native-engine/datafusion-ext-functions/src/flink_datetime.rs
new file mode 100644
index 000000000..81bbf8199
--- /dev/null
+++ b/native-engine/datafusion-ext-functions/src/flink_datetime.rs
@@ -0,0 +1,668 @@
+// Licensed to the Apache Software Foundation (ASF) under one or more
+// contributor license agreements.  See the NOTICE file distributed with
+// this work for additional information regarding copyright ownership.
+// The ASF licenses this file to You under the Apache License, Version 2.0
+// (the "License"); you may not use this file except in compliance with
+// the License.  You may obtain a copy of the License at
+//
+//    http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+use std::sync::Arc;
+
+use arrow::{
+    array::{Int64Array, StringArray},
+    datatypes::DataType,
+};
+use chrono::{DateTime, LocalResult, NaiveDateTime, Offset, TimeZone};
+use chrono_tz::{OffsetComponents, Tz};
+use datafusion::{
+    common::{DataFusionError, Result, ScalarValue},
+    physical_plan::ColumnarValue,
+};
+use datafusion_ext_commons::arrow::cast::cast;
+
+/// Native implementation of Flink SQL `UNIX_TIMESTAMP(value, format)`: parse a
+/// formatted date-time string to a Unix timestamp in seconds, replicating
+/// `java.text.SimpleDateFormat` lenient semantics.
+///
+/// Arguments are always `[value, chronoFormat, zoneId]` (arity 3). `value` is
+/// the column of strings to parse; `chronoFormat` and `zoneId` are literal
+/// scalars bound at plan time. The format is a `%`-specifier string built by
+/// the JVM-side converter from `%Y %m %d %H %M %S` plus literal characters
+/// (`%%` for a literal percent) — it is not the original Java pattern.
+///
+/// Arity 3 is the normalized form of Flink's two string-parsing arities: the
+/// JVM-side converter supplies Flink's default `yyyy-MM-dd HH:mm:ss` for the
+/// single-argument call, translates the literal pattern for the two-argument
+/// call, and resolves the session time zone at plan time. A non-literal or
+/// untranslatable format is rejected there and evaluated by Flink instead, so
+/// every call reaching this function carries all three arguments already 
bound.
+///
+/// Flink's no-argument `UNIX_TIMESTAMP()` parses nothing and yields the 
current
+/// wall-clock time per record. It is not routed here: it has no input column 
to
+/// size the output against, and it shares none of the parsing contract below.
+///
+/// An unparseable value yields `i64::MIN` (Flink's `Long.MIN_VALUE`), never
+/// NULL and never an error; a NULL value yields NULL. Arity, format and
+/// timezone problems are hard errors rather than silent defaults, so a 
plumbing
+/// bug cannot surface as silently wrong data.
+pub fn flink_unix_timestamp(args: &[ColumnarValue]) -> Result<ColumnarValue> {
+    if args.len() != 3 {
+        return Err(DataFusionError::Execution(format!(
+            "Flink_UnixTimestamp requires 3 arguments [value, chronoFormat, 
zoneId], got {}",
+            args.len()
+        )));
+    }
+
+    let format = utf8_scalar(&args[1]).ok_or_else(|| {
+        DataFusionError::Execution("Flink_UnixTimestamp: format must be a 
non-null string".into())
+    })?;
+    let zone_id = utf8_scalar(&args[2]).ok_or_else(|| {
+        DataFusionError::Execution("Flink_UnixTimestamp: zoneId must be a 
non-null string".into())
+    })?;
+    let tz: Tz = zone_id.parse().map_err(|_| {
+        DataFusionError::Execution(format!("Flink_UnixTimestamp: invalid 
timezone {zone_id}"))
+    })?;
+
+    let tokens = parse_format(&format)?;
+
+    let num_rows = match &args[0] {
+        ColumnarValue::Array(array) => array.len(),
+        ColumnarValue::Scalar(_) => 1,
+    };
+    let value = cast(&args[0].clone().into_array(num_rows)?, &DataType::Utf8)?;
+    let value = value
+        .as_any()
+        .downcast_ref::<StringArray>()
+        .expect("internal cast to Utf8 must succeed");
+
+    let result = Int64Array::from_iter(
+        value
+            .iter()
+            .map(|opt_s| opt_s.map(|s| parse_datetime(s, &tokens, 
tz).unwrap_or(i64::MIN))),
+    );
+
+    Ok(ColumnarValue::Array(Arc::new(result)))
+}
+
+fn utf8_scalar(arg: &ColumnarValue) -> Option<String> {
+    match arg {
+        ColumnarValue::Scalar(ScalarValue::Utf8(Some(s)))
+        | ColumnarValue::Scalar(ScalarValue::LargeUtf8(Some(s))) => 
Some(s.clone()),
+        _ => None,
+    }
+}
+
+/// A field is one of the six numeric SimpleDateFormat components; every
+/// supported specifier maps to exactly one. `width` is the `obeyCount` scan
+/// window used when the field is immediately followed by another numeric 
field.
+#[derive(Clone, Copy)]
+enum FieldKind {
+    Year,
+    Month,
+    Day,
+    Hour,
+    Minute,
+    Second,
+}
+
+impl FieldKind {
+    /// The `obeyCount` window width. Run length is erased by the Java→chrono
+    /// translation (both `M` and `MM` become `%m`), so a canonical width per
+    /// field is used: 4 for the year, 2 for the rest — matching the 
zero-padded
+    /// widths of Flink's default `yyyy-MM-dd HH:mm:ss`.
+    fn width(self) -> usize {
+        match self {
+            FieldKind::Year => 4,
+            _ => 2,
+        }
+    }
+}
+
+enum Token {
+    Field(FieldKind),
+    Literal(u8),
+}
+
+/// Tokenize the `%`-specifier format. Errors on an unknown specifier or a
+/// dangling `%`, since the format is a plan-time constant and such a format is
+/// a wiring bug, not bad row data.
+fn parse_format(format: &str) -> Result<Vec<Token>> {
+    let bytes = format.as_bytes();
+    let mut tokens = Vec::new();
+    let mut i = 0;
+    while i < bytes.len() {
+        if bytes[i] == b'%' {
+            let spec = bytes.get(i + 1).ok_or_else(|| {
+                DataFusionError::Execution("Flink_UnixTimestamp: dangling '%' 
in format".into())
+            })?;
+            let token = match spec {
+                b'Y' => Token::Field(FieldKind::Year),
+                b'm' => Token::Field(FieldKind::Month),
+                b'd' => Token::Field(FieldKind::Day),
+                b'H' => Token::Field(FieldKind::Hour),
+                b'M' => Token::Field(FieldKind::Minute),
+                b'S' => Token::Field(FieldKind::Second),
+                b'%' => Token::Literal(b'%'),
+                other => {
+                    return Err(DataFusionError::Execution(format!(
+                        "Flink_UnixTimestamp: unsupported format specifier 
%{}",
+                        *other as char
+                    )));
+                }
+            };
+            tokens.push(token);
+            i += 2;
+        } else {
+            tokens.push(Token::Literal(bytes[i]));
+            i += 1;
+        }
+    }
+    Ok(tokens)
+}
+
+struct Fields {
+    year: i32,
+    month: i32,
+    day: i32,
+    hour: i32,
+    minute: i32,
+    second: i32,
+}
+
+impl Default for Fields {
+    fn default() -> Self {
+        Fields {
+            year: 1970,
+            month: 1,
+            day: 1,
+            hour: 0,
+            minute: 0,
+            second: 0,
+        }
+    }
+}
+
+impl Fields {
+    fn set(&mut self, kind: FieldKind, value: i32) {
+        match kind {
+            FieldKind::Year => self.year = value,
+            FieldKind::Month => self.month = value,
+            FieldKind::Day => self.day = value,
+            FieldKind::Hour => self.hour = value,
+            FieldKind::Minute => self.minute = value,
+            FieldKind::Second => self.second = value,
+        }
+    }
+}
+
+/// Walk the tokens over `input`. Returns the Unix timestamp in seconds, or
+/// `None` on any parse failure (mapped to `i64::MIN` by the caller).
+///
+/// An empty token list consumes no input and is a failure, not a match on the
+/// epoch defaults: `DateFormat.parse` throws whenever parsing advances the
+/// position by zero characters, so Flink yields `Long.MIN_VALUE` for an empty
+/// format. Every other token consumes at least one byte — a literal matches
+/// one, a field needs at least one digit — so an empty token list is the only
+/// way to finish having consumed nothing.
+fn parse_datetime(input: &str, tokens: &[Token], tz: Tz) -> Option<i64> {
+    if tokens.is_empty() {
+        return None;
+    }
+
+    let bytes = input.as_bytes();
+    let mut pos = 0usize;
+    let mut fields = Fields::default();
+
+    for (idx, token) in tokens.iter().enumerate() {
+        match token {
+            Token::Literal(lit) => {
+                if bytes.get(pos) != Some(lit) {
+                    return None;
+                }
+                pos += 1;
+            }
+            Token::Field(kind) => {
+                // Skip ' ' and '\t' before the field, but anchor the obeyCount
+                // window at the pre-skip position so leading whitespace eats 
into
+                // the field's digit budget (a genuine SimpleDateFormat quirk).
+                let start0 = pos;
+                while matches!(bytes.get(pos), Some(b' ') | Some(b'\t')) {
+                    pos += 1;
+                }
+                if pos >= bytes.len() {
+                    return None;
+                }
+
+                // obeyCount holds when the next token is another numeric 
field;
+                // then the scan is bounded to `width` chars, otherwise it is 
greedy.
+                let obey_count = matches!(tokens.get(idx + 1), 
Some(Token::Field(_)));
+                let window_end = if obey_count {
+                    let end = start0 + kind.width();
+                    if end > bytes.len() {
+                        return None;
+                    }
+                    end
+                } else {
+                    bytes.len()
+                };
+
+                let (value, next) = scan_number(bytes, pos, window_end)?;
+                fields.set(*kind, value);
+                pos = next;
+            }
+        }
+    }
+
+    let local_sec = normalize(&fields);
+    let offset = resolve_offset_secs(local_sec, tz);
+    Some(local_sec - offset)
+}
+
+/// Scan an optional leading `-` then one or more ASCII digits within `[start,
+/// end)`. A `+` is not a sign and terminates the scan before any digit. Digits
+/// accumulate with wrapping and narrow to `i32` via low-32-bit truncation,
+/// mirroring Java's `Number.intValue()`. Returns `None` when no digit is
+/// consumed.
+fn scan_number(bytes: &[u8], start: usize, end: usize) -> Option<(i32, usize)> 
{
+    let mut pos = start;
+    let negative = bytes.get(pos) == Some(&b'-');
+    if negative {
+        pos += 1;
+    }
+
+    let digits_start = pos;
+    let mut magnitude: i64 = 0;
+    while pos < end {
+        let b = bytes[pos];
+        if !b.is_ascii_digit() {
+            break;
+        }
+        magnitude = magnitude.wrapping_mul(10).wrapping_add((b - b'0') as i64);
+        pos += 1;
+    }
+    if pos == digits_start {
+        return None;
+    }
+
+    let signed = if negative {
+        magnitude.wrapping_neg()
+    } else {
+        magnitude
+    };
+    Some((signed as i32, pos))
+}
+
+/// Normalize the (possibly out-of-range, possibly negative) fields into a 
local
+/// wall-clock time in seconds since the Unix epoch. Rollover falls out of a
+/// single month-then-day computation, so hour/minute/second overflow needs no
+/// special case. Dates before the 1582-10-15 Gregorian cutover use the Julian
+/// calendar, matching `GregorianCalendar`'s hybrid behavior (chrono is
+/// proleptic Gregorian).
+fn normalize(fields: &Fields) -> i64 {
+    let year = fields.year as i64;
+    let month = fields.month as i64;
+    let day = fields.day as i64;
+    let hour = fields.hour as i64;
+    let minute = fields.minute as i64;
+    let second = fields.second as i64;
+
+    let year_adj = year + (month - 1).div_euclid(12);
+    let month_idx = (month - 1).rem_euclid(12) + 1;
+
+    let jdn_gregorian = gregorian_jdn(year_adj, month_idx) + (day - 1);
+    let jdn = if jdn_gregorian >= 2_299_161 {
+        jdn_gregorian
+    } else {
+        julian_jdn(year_adj, month_idx) + (day - 1)
+    };
+
+    let epoch_day = jdn - 2_440_588;
+    epoch_day * 86_400 + hour * 3_600 + minute * 60 + second
+}
+
+/// Julian Day Number of the first of month `(year, month)` in the proleptic
+/// Gregorian calendar.
+fn gregorian_jdn(year: i64, month: i64) -> i64 {
+    let a = (14 - month).div_euclid(12);
+    let y = year + 4800 - a;
+    let m = month + 12 * a - 3;
+    1 + (153 * m + 2).div_euclid(5) + 365 * y + y.div_euclid(4) - 
y.div_euclid(100)
+        + y.div_euclid(400)
+        - 32045
+}
+
+/// Julian Day Number of the first of month `(year, month)` in the Julian
+/// calendar.
+fn julian_jdn(year: i64, month: i64) -> i64 {
+    let a = (14 - month).div_euclid(12);
+    let y = year + 4800 - a;
+    let m = month + 12 * a - 3;
+    1 + (153 * m + 2).div_euclid(5) + 365 * y + y.div_euclid(4) - 32083
+}
+
+/// UTC offset in seconds for a local wall-clock instant. When the local time 
is
+/// ambiguous (fall-back overlap) or nonexistent (spring-forward gap), the
+/// zone's standard (non-DST) offset is used, matching Flink.
+/// Out-of-representable-range inputs (only reachable from far-future garbage)
+/// fall back to UTC.
+fn resolve_offset_secs(local_sec: i64, tz: Tz) -> i64 {
+    let naive: NaiveDateTime = match DateTime::from_timestamp(local_sec, 0) {
+        Some(dt) => dt.naive_utc(),
+        None => return 0,
+    };
+
+    match tz.offset_from_local_datetime(&naive) {
+        LocalResult::Single(offset) => offset.fix().local_minus_utc() as i64,
+        LocalResult::Ambiguous(offset, _) => 
offset.base_utc_offset().num_seconds(),
+        LocalResult::None => tz
+            .offset_from_utc_datetime(&naive)
+            .base_utc_offset()
+            .num_seconds(),
+    }
+}
+
+#[cfg(test)]
+mod tests {
+    use arrow::array::Array;
+
+    use super::*;
+
+    /// Run one `(value, chrono_format, timezone) -> expected_seconds` case
+    /// through the full function and assert the single-row Int64 result.
+    fn run(value: &str, format: &str, tz: &str) -> i64 {
+        let args = vec![
+            ColumnarValue::Scalar(ScalarValue::Utf8(Some(value.to_string()))),
+            ColumnarValue::Scalar(ScalarValue::Utf8(Some(format.to_string()))),
+            ColumnarValue::Scalar(ScalarValue::Utf8(Some(tz.to_string()))),
+        ];
+        let out = flink_unix_timestamp(&args)
+            .expect("flink_unix_timestamp must not error on valid scalar args")
+            .into_array(1)
+            .expect("result must materialize to an array");
+        let out = out
+            .as_any()
+            .downcast_ref::<Int64Array>()
+            .expect("result must be Int64Array");
+        assert_eq!(out.len(), 1);
+        out.value(0)
+    }
+
+    const MIN: i64 = i64::MIN;
+    const DEFAULT: &str = "%Y-%m-%d %H:%M:%S";
+
+    // Happy path across formats and timezones.
+    #[test]
+    fn oracle_happy_path() {
+        assert_eq!(run("2020-10-10 00:00:01", DEFAULT, "UTC"), 1602288001);
+        assert_eq!(
+            run("2020-10-10 00:00:01", DEFAULT, "Asia/Shanghai"),
+            1602259201
+        );
+        assert_eq!(
+            run("2015-07-24 10:00:00", DEFAULT, "Asia/Tokyo"),
+            1437699600
+        );
+        assert_eq!(run("20201010000001", "%Y%m%d%H%M%S", "UTC"), 1602288001);
+        assert_eq!(run("2020-10-10", "%Y-%m-%d", "UTC"), 1602288000);
+    }
+
+    // Non-padded fields (1-digit values under 2-digit specifiers) — the core
+    // leniency requirement.
+    #[test]
+    fn oracle_non_padded() {
+        assert_eq!(run("2020-1-1 0:0:0", DEFAULT, "UTC"), 1577836800);
+        assert_eq!(run("2020-01-01 00:00:00", DEFAULT, "UTC"), 1577836800);
+    }
+
+    // Greedy scan (field bounded only by the next literal) and obeyCount 
windows
+    // (field bounded to its canonical width when followed by another numeric
+    // field).
+    #[test]
+    fn oracle_greedy_and_obey_count() {
+        assert_eq!(run("2020-115-01 00:00:00", DEFAULT, "UTC"), 1877558400);
+        assert_eq!(run("12020-01-01 00:00:00", DEFAULT, "UTC"), 317147356800);
+        assert_eq!(run("2020-0000001-01 00:00:00", DEFAULT, "UTC"), 
1577836800);
+        assert_eq!(run("2020101000000", "%Y%m%d%H%M%S", "UTC"), 1602288000);
+        assert_eq!(run("202010100000012", "%Y%m%d%H%M%S", "UTC"), 1602288012);
+    }
+
+    // Field rollover across month/day/hour/minute/second, including the
+    // normalization order that rolls month before day.
+    #[test]
+    fn oracle_rollover() {
+        assert_eq!(run("2020-13-01 00:00:00", DEFAULT, "UTC"), 1609459200);
+        assert_eq!(run("2020-00-01 00:00:00", DEFAULT, "UTC"), 1575158400);
+        assert_eq!(run("2020-99-01 00:00:00", DEFAULT, "UTC"), 1835481600);
+        assert_eq!(run("2020-01-45 00:00:00", DEFAULT, "UTC"), 1581638400);
+        assert_eq!(run("2020-01-00 00:00:00", DEFAULT, "UTC"), 1577750400);
+        assert_eq!(run("2020-02-30 00:00:00", DEFAULT, "UTC"), 1583020800);
+        assert_eq!(run("2021-02-30 00:00:00", DEFAULT, "UTC"), 1614643200);
+        assert_eq!(run("2020-01-01 24:00:00", DEFAULT, "UTC"), 1577923200);
+        assert_eq!(run("2020-01-01 25:00:00", DEFAULT, "UTC"), 1577926800);
+        assert_eq!(run("2020-01-01 00:60:00", DEFAULT, "UTC"), 1577840400);
+        assert_eq!(run("2020-01-01 00:00:60", DEFAULT, "UTC"), 1577836860);
+        assert_eq!(run("2020-01-01 00:00:99", DEFAULT, "UTC"), 1577836899);
+        assert_eq!(run("2020-13-32 00:00:00", DEFAULT, "UTC"), 1612137600);
+        assert_eq!(run("2020-25-32 00:00:00", DEFAULT, "UTC"), 1643673600);
+        assert_eq!(run("2020-13-32 25:61:61", DEFAULT, "UTC"), 1612231321);
+    }
+
+    // Negative fields: a leading '-' is a valid sign (negative month, 
negative day,
+    // BC year through the Julian branch).
+    #[test]
+    fn oracle_negative_fields() {
+        assert_eq!(run("2020--5-10 00:00:01", DEFAULT, "UTC"), 1562716801);
+        assert_eq!(run("2020-10--5 00:00:01", DEFAULT, "UTC"), 1600992001);
+        assert_eq!(run("-2020-10-10 00:00:01", DEFAULT, "UTC"), -125889292799);
+    }
+
+    // Trailing input is ignored; leading/inner whitespace is skipped before a 
field
+    // (and eats the obeyCount budget) but never before a literal.
+    #[test]
+    fn oracle_whitespace_and_trailing() {
+        assert_eq!(run("2020-10-10 00:00:01XYZ", DEFAULT, "UTC"), 1602288001);
+        assert_eq!(run("2020-10-10 00:00:01   ", DEFAULT, "UTC"), 1602288001);
+        assert_eq!(run("2020-10-10 00:00:01", "%Y-%m-%d", "UTC"), 1602288000);
+        assert_eq!(run(" 2020-10-10 00:00:01", DEFAULT, "UTC"), 1602288001);
+        assert_eq!(run("\t2020-10-10 00:00:01", DEFAULT, "UTC"), 1602288001);
+        assert_eq!(run("X2020-10-10 00:00:01", DEFAULT, "UTC"), MIN);
+        assert_eq!(run("2020- 10-10 00:00:01", DEFAULT, "UTC"), 1602288001);
+        assert_eq!(run("2020 -10-10 00:00:01", DEFAULT, "UTC"), MIN);
+        assert_eq!(run(" 20200101", "%Y%m%d", "UTC"), -55786752000);
+    }
+
+    // Parse failures resolve to i64::MIN, never NULL and never an error: 
literal
+    // mismatch, short input, non-digit field, grouping separator, and a '+' 
sign.
+    #[test]
+    fn oracle_failures() {
+        assert_eq!(run("not a date", DEFAULT, "UTC"), MIN);
+        assert_eq!(run("", DEFAULT, "UTC"), MIN);
+        assert_eq!(run("   ", DEFAULT, "UTC"), MIN);
+        assert_eq!(run("2020-10-10", DEFAULT, "UTC"), MIN);
+        assert_eq!(run("2020/10/10 00:00:01", DEFAULT, "UTC"), MIN);
+        assert_eq!(run("2020-XX-10 00:00:01", DEFAULT, "UTC"), MIN);
+        assert_eq!(run("2,020-01-01 00:00:01", DEFAULT, "UTC"), MIN);
+        assert_eq!(run("+2020-10-10 00:00:01", DEFAULT, "UTC"), MIN);
+    }
+
+    // Epoch and pre-1970 negatives.
+    //
+    // Sub-second inputs (e.g. `... .SSS` with millisecond truncation) are not
+    // tested here: the supported specifiers carry no sub-second precision, and
+    // a format containing `S` is rejected by the upstream scanner, so such
+    // inputs fall back to Flink and never reach this function.
+    #[test]
+    fn oracle_epoch() {
+        assert_eq!(run("1970-01-01 00:00:00", DEFAULT, "UTC"), 0);
+        assert_eq!(run("1969-12-31 23:59:59", DEFAULT, "UTC"), -1);
+        assert_eq!(run("1900-01-01 00:00:00", DEFAULT, "UTC"), -2208988800);
+    }
+
+    // DST gap and overlap across several zones — the zone's standard (non-DST)
+    // offset wins, so a spring-forward gap collapses to the post-transition
+    // instant.
+    #[test]
+    fn oracle_dst() {
+        assert_eq!(run("2005-03-27 02:00:00", DEFAULT, "MET"), 1111885200);
+        assert_eq!(run("2005-03-27 03:00:00", DEFAULT, "MET"), 1111885200);
+        assert_eq!(
+            run("2021-03-14 02:30:00", DEFAULT, "America/New_York"),
+            1615707000
+        );
+        assert_eq!(run("2005-10-30 02:30:00", DEFAULT, "MET"), 1130635800);
+        assert_eq!(
+            run("2021-11-07 01:30:00", DEFAULT, "America/New_York"),
+            1636266600
+        );
+        assert_eq!(
+            run("2021-04-04 02:30:00", DEFAULT, "Australia/Sydney"),
+            1617467400
+        );
+        assert_eq!(
+            run("2021-10-03 02:30:00", DEFAULT, "Australia/Sydney"),
+            1633192200
+        );
+    }
+
+    // Literals from quoted patterns: a case-sensitive letter, a literal 
percent
+    // (`%%`), and a literal single quote.
+    #[test]
+    fn oracle_quoting() {
+        assert_eq!(
+            run("2020-10-10T00:00:01", "%Y-%m-%dT%H:%M:%S", "UTC"),
+            1602288001
+        );
+        assert_eq!(run("2020-10-10t00:00:01", "%Y-%m-%dT%H:%M:%S", "UTC"), 
MIN);
+        assert_eq!(run("2020%10", "%Y%%%m", "UTC"), 1601510400);
+        assert_eq!(run("2020'10", "%Y'%m", "UTC"), 1601510400);
+    }
+
+    // Julian/Gregorian hybrid calendar around the 1582-10-15 cutover: the last
+    // Julian day, the first Gregorian day, a nonexistent day mapped forward, 
and a
+    // deep-Julian date.
+    #[test]
+    fn oracle_hybrid_calendar() {
+        assert_eq!(run("1582-10-15 00:00:00", DEFAULT, "UTC"), -12219292800);
+        assert_eq!(run("1582-10-04 00:00:00", DEFAULT, "UTC"), -12219379200);
+        assert_eq!(run("1582-10-05 00:00:00", DEFAULT, "UTC"), -12219292800);
+        assert_eq!(run("1000-01-01 00:00:00", DEFAULT, "UTC"), -30609792000);
+    }
+
+    // A NULL value yields NULL output, distinct from the i64::MIN 
parse-failure
+    // sentinel.
+    #[test]
+    fn null_value_yields_null() {
+        let args = vec![
+            ColumnarValue::Scalar(ScalarValue::Utf8(None)),
+            
ColumnarValue::Scalar(ScalarValue::Utf8(Some(DEFAULT.to_string()))),
+            ColumnarValue::Scalar(ScalarValue::Utf8(Some("UTC".to_string()))),
+        ];
+        let out = flink_unix_timestamp(&args)
+            .expect("must not error")
+            .into_array(1)
+            .expect("must materialize");
+        let out = out
+            .as_any()
+            .downcast_ref::<Int64Array>()
+            .expect("result must be Int64Array");
+        assert_eq!(out.len(), 1);
+        assert!(out.is_null(0));
+    }
+
+    // Mixed array: NULL propagates per-row, valid rows parse, garbage → 
i64::MIN.
+    #[test]
+    fn array_with_null_and_garbage() {
+        let value = Arc::new(StringArray::from(vec![
+            Some("2020-10-10 00:00:01"),
+            None,
+            Some("garbage"),
+        ]));
+        let args = vec![
+            ColumnarValue::Array(value),
+            
ColumnarValue::Scalar(ScalarValue::Utf8(Some(DEFAULT.to_string()))),
+            ColumnarValue::Scalar(ScalarValue::Utf8(Some("UTC".to_string()))),
+        ];
+        let out = flink_unix_timestamp(&args)
+            .expect("must not error")
+            .into_array(3)
+            .expect("must materialize");
+        let out = out
+            .as_any()
+            .downcast_ref::<Int64Array>()
+            .expect("result must be Int64Array");
+        assert_eq!(out.value(0), 1602288001);
+        assert!(out.is_null(1));
+        assert_eq!(out.value(2), i64::MIN);
+    }
+
+    // In numeric adjacency, `%m` scans its canonical 2-digit obeyCount 
window: on
+    // `20201010` with `%Y%m%d`, `%Y` reads 4 (2020), `%m` reads 2 (10), `%d` 
reads
+    // the rest (10), giving 2020-10-10. The run length that would distinguish 
a
+    // 1-digit month is not carried in the chrono format, so the native 
canonical
+    // width defines the contract here (the JVM scanner decides which patterns 
reach
+    // native).
+    #[test]
+    fn adjacent_month_uses_canonical_width_two() {
+        assert_eq!(run("20201010", "%Y%m%d", "UTC"), 1602288000);
+    }
+
+    // The year field's canonical obeyCount window is 4: on `20200101` with
+    // `%Y%m%d`, `%Y` reads exactly `2020`, leaving `01`/`01` for month and 
day.
+    #[test]
+    fn year_uses_canonical_width_four() {
+        assert_eq!(run("20200101", "%Y%m%d", "UTC"), 1577836800);
+    }
+
+    #[test]
+    fn arity_mismatch_errors() {
+        let two = vec![
+            ColumnarValue::Scalar(ScalarValue::Utf8(Some("2020-10-10 
00:00:01".to_string()))),
+            
ColumnarValue::Scalar(ScalarValue::Utf8(Some(DEFAULT.to_string()))),
+        ];
+        assert!(flink_unix_timestamp(&two).is_err());
+
+        let four = vec![
+            ColumnarValue::Scalar(ScalarValue::Utf8(Some("2020-10-10 
00:00:01".to_string()))),
+            
ColumnarValue::Scalar(ScalarValue::Utf8(Some(DEFAULT.to_string()))),
+            ColumnarValue::Scalar(ScalarValue::Utf8(Some("UTC".to_string()))),
+            
ColumnarValue::Scalar(ScalarValue::Utf8(Some("extra".to_string()))),
+        ];
+        assert!(flink_unix_timestamp(&four).is_err());
+    }
+
+    #[test]
+    fn invalid_timezone_errors() {
+        let args = vec![
+            ColumnarValue::Scalar(ScalarValue::Utf8(Some("2020-10-10 
00:00:01".to_string()))),
+            
ColumnarValue::Scalar(ScalarValue::Utf8(Some(DEFAULT.to_string()))),
+            
ColumnarValue::Scalar(ScalarValue::Utf8(Some("Mars/Olympus".to_string()))),
+        ];
+        assert!(flink_unix_timestamp(&args).is_err());
+    }
+
+    // An empty format consumes nothing, so it matches no input at all rather
+    // than matching everything at the epoch defaults.
+    #[test]
+    fn empty_format_never_matches() {
+        assert_eq!(run("2020-10-10 00:00:01", "", "UTC"), MIN);
+        assert_eq!(run("2020-10-10 00:00:01", "", "Asia/Shanghai"), MIN);
+        assert_eq!(run("", "", "UTC"), MIN);
+    }
+
+    #[test]
+    fn unsupported_specifier_errors() {
+        let args = vec![
+            ColumnarValue::Scalar(ScalarValue::Utf8(Some("2020".to_string()))),
+            
ColumnarValue::Scalar(ScalarValue::Utf8(Some("%Y-%Q".to_string()))),
+            ColumnarValue::Scalar(ScalarValue::Utf8(Some("UTC".to_string()))),
+        ];
+        assert!(flink_unix_timestamp(&args).is_err());
+    }
+}
diff --git a/native-engine/datafusion-ext-functions/src/lib.rs 
b/native-engine/datafusion-ext-functions/src/lib.rs
index 7d90154db..d7fecc982 100644
--- a/native-engine/datafusion-ext-functions/src/lib.rs
+++ b/native-engine/datafusion-ext-functions/src/lib.rs
@@ -19,6 +19,7 @@ use datafusion::{common::Result, 
logical_expr::ScalarFunctionImplementation};
 use datafusion_ext_commons::df_unimplemented_err;
 
 mod brickhouse;
+mod flink_datetime;
 mod spark_array;
 mod spark_bround;
 mod spark_check_overflow;
@@ -97,6 +98,7 @@ pub fn create_auron_ext_function(
             
Arc::new(spark_normalize_nan_and_zero::spark_normalize_nan_and_zero)
         }
         "Spark_IsNaN" => Arc::new(spark_isnan::spark_isnan),
-        _ => df_unimplemented_err!("spark ext function not implemented: 
{name}")?,
+        "Flink_UnixTimestamp" => 
Arc::new(flink_datetime::flink_unix_timestamp),
+        _ => df_unimplemented_err!("auron ext function not implemented: 
{name}")?,
     })
 }

Reply via email to