@@ -52,7 +52,10 @@ use skardi::sources::providers::{
5252 lance:: register_lance_table,
5353 mongo:: register_mongo_tables,
5454 mysql:: register_mysql_tables,
55- open_connector:: { OpenConnectorConfig , register_open_connector_tables} ,
55+ open_connector:: {
56+ OpenConnectorConfig , OpenConnectorGateways , register_open_connector_tables,
57+ register_open_connector_udtfs,
58+ } ,
5659 sqlite:: {
5760 register_sqlite_fts_udtf, register_sqlite_knn_udtf, register_sqlite_tables,
5861 register_vec_to_binary_udf,
@@ -479,9 +482,10 @@ impl UrlTableFactory for SkardiUrlTableFactory {
479482}
480483
481484/// Create a new SessionContext with custom URL table support (built-in files + Lance)
482- /// and the `lance_knn` / `pg_knn` UDTFs registered.
483- fn new_session_context ( ) -> ( SessionContext , DatasetRegistry ) {
485+ /// and the `lance_knn` / `pg_knn` / Open Connector UDTFs registered.
486+ fn new_session_context ( ) -> ( SessionContext , DatasetRegistry , OpenConnectorGateways ) {
484487 let dataset_registry: DatasetRegistry = Arc :: new ( RwLock :: new ( HashMap :: new ( ) ) ) ;
488+ let open_connector_gateways = OpenConnectorGateways :: default ( ) ;
485489 let session_store = SessionStore :: new ( ) ;
486490 let factory = Arc :: new ( SkardiUrlTableFactory :: new (
487491 session_store,
@@ -516,6 +520,9 @@ fn new_session_context() -> (SessionContext, DatasetRegistry) {
516520 register_sqlite_knn_udtf ( & ctx, Arc :: clone ( & dataset_registry) ) ;
517521 register_sqlite_fts_udtf ( & ctx, Arc :: clone ( & dataset_registry) ) ;
518522 register_vec_to_binary_udf ( & mut ctx) ;
523+ // Open Connector UDTFs plan against the gateway state that
524+ // register_open_connector_tables fills in during ctx registration.
525+ register_open_connector_udtfs ( & ctx, Arc :: clone ( & open_connector_gateways) ) ;
519526
520527 // Embedding UDFs (gated by feature flags, lazy model loading on first call).
521528 #[ cfg( feature = "onnx" ) ]
@@ -544,7 +551,7 @@ fn new_session_context() -> (SessionContext, DatasetRegistry) {
544551 registry. register_chunk_udf ( & mut ctx) ;
545552 }
546553
547- ( ctx, dataset_registry)
554+ ( ctx, dataset_registry, open_connector_gateways )
548555}
549556
550557/// Resolve a path string: if relative (and not remote), resolve against cwd.
@@ -733,13 +740,19 @@ async fn load_and_register_all(
733740 ctx_path : & Path ,
734741 session_ctx : & mut SessionContext ,
735742 dataset_registry : & DatasetRegistry ,
743+ open_connector_gateways : & OpenConnectorGateways ,
736744) -> Result < LocalContextConfig > {
737745 let config = read_context_file ( ctx_path) ?;
738746
739747 for source in & config. data_sources {
740- register_source ( session_ctx, source, dataset_registry)
741- . await
742- . with_context ( || format ! ( "Failed to register data source '{}'" , source. name) ) ?;
748+ register_source (
749+ session_ctx,
750+ source,
751+ dataset_registry,
752+ open_connector_gateways,
753+ )
754+ . await
755+ . with_context ( || format ! ( "Failed to register data source '{}'" , source. name) ) ?;
743756 }
744757
745758 Ok ( config)
@@ -774,6 +787,7 @@ async fn register_source(
774787 session_ctx : & mut SessionContext ,
775788 source : & LocalDataSource ,
776789 dataset_registry : & DatasetRegistry ,
790+ open_connector_gateways : & OpenConnectorGateways ,
777791) -> Result < ( ) > {
778792 let source_type = source. source_type . to_lowercase ( ) ;
779793
@@ -957,6 +971,7 @@ async fn register_source(
957971 source. open_connector . as_ref ( ) ,
958972 source. is_read_write ( ) ,
959973 source. hierarchy_level ,
974+ Some ( open_connector_gateways) ,
960975 )
961976 . await
962977 . with_context ( || format ! ( "Failed to register Open Connector '{}'" , source. name) ) ?;
@@ -1222,8 +1237,9 @@ async fn show_schema(
12221237 table_filter : Option < & str > ,
12231238 out : & mut dyn Write ,
12241239) -> Result < ( ) > {
1225- let ( mut session_ctx, dataset_registry) = new_session_context ( ) ;
1226- let config = load_and_register_all ( ctx_path, & mut session_ctx, & dataset_registry) . await ?;
1240+ let ( mut session_ctx, dataset_registry, oc_gateways) = new_session_context ( ) ;
1241+ let config =
1242+ load_and_register_all ( ctx_path, & mut session_ctx, & dataset_registry, & oc_gateways) . await ?;
12271243
12281244 // Build the semantics registry from the same inputs the server uses:
12291245 // - the ctx-inline `description` field on each data source (fallback), and
@@ -1368,12 +1384,13 @@ fn source_name_for<'a>(
13681384/// data sources first. If no context file is found, run the query in a bare session with
13691385/// URL table support (allowing direct file/lance paths in SQL).
13701386async fn run_query ( ctx_override : Option < PathBuf > , sql : & str ) -> Result < ( ) > {
1371- let ( mut session_ctx, dataset_registry) = new_session_context ( ) ;
1387+ let ( mut session_ctx, dataset_registry, oc_gateways ) = new_session_context ( ) ;
13721388
13731389 // Try to load context file, but don't fail if not found when no explicit --ctx was given
13741390 match resolve_ctx_path ( ctx_override. as_deref ( ) ) {
13751391 Some ( ctx_path) if ctx_path. exists ( ) => {
1376- load_and_register_all ( & ctx_path, & mut session_ctx, & dataset_registry) . await ?;
1392+ load_and_register_all ( & ctx_path, & mut session_ctx, & dataset_registry, & oc_gateways)
1393+ . await ?;
13771394 }
13781395 Some ( ctx_path) if ctx_override. is_some ( ) => {
13791396 anyhow:: bail!( "Context file not found: {}" , ctx_path. display( ) ) ;
@@ -1525,9 +1542,9 @@ async fn run_pipeline_with_params(
15251542 . with_context ( || format ! ( "Failed to render SQL for pipeline '{}'" , pipeline_name) ) ?;
15261543
15271544 // 3. Build a SessionContext with ctx data sources registered, then execute.
1528- let ( mut session_ctx, dataset_registry) = new_session_context ( ) ;
1545+ let ( mut session_ctx, dataset_registry, oc_gateways ) = new_session_context ( ) ;
15291546 if let Some ( p) = & ctx_path_for_load {
1530- load_and_register_all ( p, & mut session_ctx, & dataset_registry) . await ?;
1547+ load_and_register_all ( p, & mut session_ctx, & dataset_registry, & oc_gateways ) . await ?;
15311548 }
15321549
15331550 auto_register_object_stores_from_sql ( & session_ctx, & sql) ?;
@@ -2879,10 +2896,15 @@ spec:
28792896
28802897 #[ tokio:: test]
28812898 async fn errors_without_connection_string ( ) {
2882- let ( mut session_ctx, registry) = new_session_context ( ) ;
2883- let err = register_source ( & mut session_ctx, & dynamodb_source ( None ) , & registry)
2884- . await
2885- . unwrap_err ( ) ;
2899+ let ( mut session_ctx, registry, oc_gateways) = new_session_context ( ) ;
2900+ let err = register_source (
2901+ & mut session_ctx,
2902+ & dynamodb_source ( None ) ,
2903+ & registry,
2904+ & oc_gateways,
2905+ )
2906+ . await
2907+ . unwrap_err ( ) ;
28862908 let msg = format ! ( "{err:?}" ) ;
28872909 assert ! (
28882910 msg. contains( "connection_string (endpoint URL) required" ) ,
@@ -2892,9 +2914,9 @@ spec:
28922914
28932915 #[ tokio:: test]
28942916 async fn errors_without_options ( ) {
2895- let ( mut session_ctx, registry) = new_session_context ( ) ;
2917+ let ( mut session_ctx, registry, oc_gateways ) = new_session_context ( ) ;
28962918 let source = dynamodb_source ( Some ( "http://localhost:8000" ) ) ;
2897- let err = register_source ( & mut session_ctx, & source, & registry)
2919+ let err = register_source ( & mut session_ctx, & source, & registry, & oc_gateways )
28982920 . await
28992921 . unwrap_err ( ) ;
29002922 let msg = format ! ( "{err:?}" ) ;
@@ -2927,10 +2949,15 @@ spec:
29272949
29282950 #[ tokio:: test]
29292951 async fn errors_without_connection_string ( ) {
2930- let ( mut session_ctx, registry) = new_session_context ( ) ;
2931- let err = register_source ( & mut session_ctx, & clickhouse_source ( None ) , & registry)
2932- . await
2933- . unwrap_err ( ) ;
2952+ let ( mut session_ctx, registry, oc_gateways) = new_session_context ( ) ;
2953+ let err = register_source (
2954+ & mut session_ctx,
2955+ & clickhouse_source ( None ) ,
2956+ & registry,
2957+ & oc_gateways,
2958+ )
2959+ . await
2960+ . unwrap_err ( ) ;
29342961 let msg = format ! ( "{err:?}" ) ;
29352962 assert ! (
29362963 msg. contains( "connection_string required" ) ,
@@ -2943,10 +2970,10 @@ spec:
29432970 // The provider is the single enforcement point for the read-only
29442971 // invariant — the CLI must reject read_write exactly like the
29452972 // server's UnsupportedWriteMode.
2946- let ( mut session_ctx, registry) = new_session_context ( ) ;
2973+ let ( mut session_ctx, registry, oc_gateways ) = new_session_context ( ) ;
29472974 let mut source = clickhouse_source ( Some ( "http://127.0.0.1:1" ) ) ;
29482975 source. access_mode = Some ( "read_write" . to_string ( ) ) ;
2949- let err = register_source ( & mut session_ctx, & source, & registry)
2976+ let err = register_source ( & mut session_ctx, & source, & registry, & oc_gateways )
29502977 . await
29512978 . unwrap_err ( ) ;
29522979 let msg = format ! ( "{err:?}" ) ;
@@ -2988,9 +3015,9 @@ bindings:
29883015
29893016 #[ tokio:: test]
29903017 async fn errors_without_connection_string ( ) {
2991- let ( mut session_ctx, registry) = new_session_context ( ) ;
3018+ let ( mut session_ctx, registry, oc_gateways ) = new_session_context ( ) ;
29923019 let source = open_connector_source ( None , Some ( VALID_CONFIG ) ) ;
2993- let err = register_source ( & mut session_ctx, & source, & registry)
3020+ let err = register_source ( & mut session_ctx, & source, & registry, & oc_gateways )
29943021 . await
29953022 . unwrap_err ( ) ;
29963023 let msg = format ! ( "{err:?}" ) ;
@@ -3004,11 +3031,11 @@ bindings:
30043031 async fn errors_with_table_hierarchy ( ) {
30053032 // hierarchy_level defaults to Table; the CLI must reject it with
30063033 // a clear message, not the provider's wrapped error.
3007- let ( mut session_ctx, registry) = new_session_context ( ) ;
3034+ let ( mut session_ctx, registry, oc_gateways ) = new_session_context ( ) ;
30083035 let mut source =
30093036 open_connector_source ( Some ( "http://localhost:3000" ) , Some ( VALID_CONFIG ) ) ;
30103037 source. hierarchy_level = HierarchyLevel :: Table ;
3011- let err = register_source ( & mut session_ctx, & source, & registry)
3038+ let err = register_source ( & mut session_ctx, & source, & registry, & oc_gateways )
30123039 . await
30133040 . unwrap_err ( ) ;
30143041 let msg = format ! ( "{err:?}" ) ;
@@ -3020,9 +3047,9 @@ bindings:
30203047
30213048 #[ tokio:: test]
30223049 async fn errors_without_typed_config ( ) {
3023- let ( mut session_ctx, registry) = new_session_context ( ) ;
3050+ let ( mut session_ctx, registry, oc_gateways ) = new_session_context ( ) ;
30243051 let source = open_connector_source ( Some ( "http://localhost:3000" ) , None ) ;
3025- let err = register_source ( & mut session_ctx, & source, & registry)
3052+ let err = register_source ( & mut session_ctx, & source, & registry, & oc_gateways )
30263053 . await
30273054 . unwrap_err ( ) ;
30283055 let msg = format ! ( "{err:?}" ) ;
@@ -3038,11 +3065,11 @@ bindings:
30383065 // The provider is the single enforcement point for the
30393066 // read-only invariant — the CLI must reject read_write exactly
30403067 // like the server's UnsupportedWriteMode.
3041- let ( mut session_ctx, registry) = new_session_context ( ) ;
3068+ let ( mut session_ctx, registry, oc_gateways ) = new_session_context ( ) ;
30423069 let mut source =
30433070 open_connector_source ( Some ( "http://localhost:3000" ) , Some ( VALID_CONFIG ) ) ;
30443071 source. access_mode = Some ( "read_write" . to_string ( ) ) ;
3045- let err = register_source ( & mut session_ctx, & source, & registry)
3072+ let err = register_source ( & mut session_ctx, & source, & registry, & oc_gateways )
30463073 . await
30473074 . unwrap_err ( ) ;
30483075 let msg = format ! ( "{err:?}" ) ;
@@ -3051,11 +3078,11 @@ bindings:
30513078
30523079 #[ tokio:: test]
30533080 async fn errors_when_typed_config_on_wrong_type ( ) {
3054- let ( mut session_ctx, registry) = new_session_context ( ) ;
3081+ let ( mut session_ctx, registry, oc_gateways ) = new_session_context ( ) ;
30553082 let mut source =
30563083 open_connector_source ( Some ( "http://localhost:3000" ) , Some ( VALID_CONFIG ) ) ;
30573084 source. source_type = "csv" . to_string ( ) ;
3058- let err = register_source ( & mut session_ctx, & source, & registry)
3085+ let err = register_source ( & mut session_ctx, & source, & registry, & oc_gateways )
30593086 . await
30603087 . unwrap_err ( ) ;
30613088 let msg = format ! ( "{err:?}" ) ;
@@ -3069,11 +3096,11 @@ bindings:
30693096 async fn errors_when_token_env_missing ( ) {
30703097 // With the config valid, the next failure is the unset runtime
30713098 // token — before any network call to the (unroutable) gateway.
3072- let ( mut session_ctx, registry) = new_session_context ( ) ;
3099+ let ( mut session_ctx, registry, oc_gateways ) = new_session_context ( ) ;
30733100 let config =
30743101 VALID_CONFIG . replace ( "OPEN_CONNECTOR_TOKEN" , "SKARDI_CLI_TEST_OC_TOKEN_UNSET" ) ;
30753102 let source = open_connector_source ( Some ( "http://127.0.0.1:1" ) , Some ( config. as_str ( ) ) ) ;
3076- let err = register_source ( & mut session_ctx, & source, & registry)
3103+ let err = register_source ( & mut session_ctx, & source, & registry, & oc_gateways )
30773104 . await
30783105 . unwrap_err ( ) ;
30793106 let msg = format ! ( "{err:?}" ) ;
0 commit comments