Skip to content

Commit 00da0f1

Browse files
davidwzhaoclaude
andauthored
Add storage integration to CSVConfig (#257)
* Add storage integration to CSVConfig * Add storage_integration to CSV grammar Make CSVConfig.storage_integration optional and parse/print it as a nested (storage_integration {...}) config_dict. Secret credentials are masked as "***". Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * Add storage_integration tests; fix printer name collision Rename the grammar nonterm to storage_integration so it no longer collides with the auto-generated CSVStorageIntegration message printer, and regenerate the SDK parsers/printers. Add a csv_storage_integration fixture (S3/Azure import + S3 export) plus unit tests covering secret masking. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * Rename CSVStorageIntegration to StorageIntegration Rename the grammar nonterm to _storage_integration to avoid colliding with the renamed message's printer, and regenerate the SDKs. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
1 parent d6a3a89 commit 00da0f1

17 files changed

Lines changed: 14835 additions & 13990 deletions

File tree

meta/src/meta/grammar.y

Lines changed: 53 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -91,6 +91,7 @@
9191
%nonterm csv_data logic.CSVData
9292
%nonterm csv_locator_inline_data String
9393
%nonterm csv_locator_paths Sequence[String]
94+
%nonterm _storage_integration Sequence[Tuple[String, logic.Value]]
9495
%nonterm csvlocator logic.CSVLocator
9596
%nonterm data logic.Data
9697
%nonterm date logic.DateValue
@@ -216,6 +217,7 @@
216217
%validator_ignore_completeness BeTreeLocator
217218
%validator_ignore_completeness BeTreeConfig
218219
%validator_ignore_completeness ExportCSVColumns
220+
%validator_ignore_completeness StorageIntegration
219221

220222
%%
221223

@@ -1127,9 +1129,16 @@ csvlocator
11271129
$4: Optional[String] = builtin.decode_string($$.inline_data) if builtin.decode_string($$.inline_data) != "" else None
11281130

11291131
csv_config
1130-
: "(" "csv_config" config_dict ")"
1131-
construct: $$ = construct_csv_config($3)
1132-
deconstruct: $3: Sequence[Tuple[String, logic.Value]] = deconstruct_csv_config($$)
1132+
: "(" "csv_config" config_dict _storage_integration? ")"
1133+
construct: $$ = construct_csv_config($3, $4)
1134+
deconstruct:
1135+
$3: Sequence[Tuple[String, logic.Value]] = deconstruct_csv_config($$)
1136+
$4: Optional[Sequence[Tuple[String, logic.Value]]] = deconstruct_csv_storage_integration_optional($$)
1137+
1138+
_storage_integration
1139+
: "(" "storage_integration" config_dict ")"
1140+
construct: $$ = $3
1141+
deconstruct: $3: Sequence[Tuple[String, logic.Value]] = $$
11331142

11341143
gnf_column_path
11351144
: STRING
@@ -1464,7 +1473,10 @@ def _try_extract_value_string_list(value: Optional[logic.Value]) -> Optional[Seq
14641473
return None
14651474

14661475

1467-
def construct_csv_config(config_dict: Sequence[Tuple[String, logic.Value]]) -> logic.CSVConfig:
1476+
def construct_csv_config(
1477+
config_dict: Sequence[Tuple[String, logic.Value]],
1478+
storage_integration_opt: Optional[Sequence[Tuple[String, logic.Value]]],
1479+
) -> logic.CSVConfig:
14681480
config: Dict[String, logic.Value] = builtin.dict_from_list(config_dict)
14691481
header_row: Int32 = _extract_value_int32(builtin.dict_get(config, "csv_header_row"), 1)
14701482
skip: int = _extract_value_int64(builtin.dict_get(config, "csv_skip"), 0)
@@ -1478,6 +1490,7 @@ def construct_csv_config(config_dict: Sequence[Tuple[String, logic.Value]]) -> l
14781490
encoding: str = _extract_value_string(builtin.dict_get(config, "csv_encoding"), "utf-8")
14791491
compression: str = _extract_value_string(builtin.dict_get(config, "csv_compression"), "auto")
14801492
partition_size_mb: int = _extract_value_int64(builtin.dict_get(config, "csv_partition_size_mb"), 0)
1493+
storage_integration: Optional[logic.StorageIntegration] = construct_csv_storage_integration(storage_integration_opt)
14811494
return logic.CSVConfig(
14821495
header_row=header_row,
14831496
skip=skip,
@@ -1491,9 +1504,25 @@ def construct_csv_config(config_dict: Sequence[Tuple[String, logic.Value]]) -> l
14911504
encoding=encoding,
14921505
compression=compression,
14931506
partition_size_mb=partition_size_mb,
1507+
storage_integration=storage_integration,
14941508
)
14951509

14961510

1511+
def construct_csv_storage_integration(
1512+
storage_integration_opt: Optional[Sequence[Tuple[String, logic.Value]]],
1513+
) -> Optional[logic.StorageIntegration]:
1514+
if storage_integration_opt is None:
1515+
return builtin.none()
1516+
config: Dict[String, logic.Value] = builtin.dict_from_list(builtin.unwrap_option(storage_integration_opt))
1517+
return builtin.some(logic.StorageIntegration(
1518+
provider=_extract_value_string(builtin.dict_get(config, "provider"), ""),
1519+
azure_sas_token=_extract_value_string(builtin.dict_get(config, "azure_sas_token"), ""),
1520+
s3_region=_extract_value_string(builtin.dict_get(config, "s3_region"), ""),
1521+
s3_access_key_id=_extract_value_string(builtin.dict_get(config, "s3_access_key_id"), ""),
1522+
s3_secret_access_key=_extract_value_string(builtin.dict_get(config, "s3_secret_access_key"), ""),
1523+
))
1524+
1525+
14971526
def construct_betree_info(
14981527
key_types: Sequence[logic.Type],
14991528
value_types: Sequence[logic.Type],
@@ -1658,6 +1687,26 @@ def deconstruct_csv_config(msg: logic.CSVConfig) -> List[Tuple[String, logic.Val
16581687
return builtin.list_sort(result)
16591688

16601689

1690+
# Secret credential values are masked in the human-readable output. As a result the
1691+
# storage integration block does not round-trip the real secrets back through the parser.
1692+
def deconstruct_csv_storage_integration_optional(msg: logic.CSVConfig) -> Optional[Sequence[Tuple[String, logic.Value]]]:
1693+
if not builtin.has_proto_field(msg, "storage_integration"):
1694+
return builtin.none()
1695+
si: logic.StorageIntegration = builtin.unwrap_option(msg.storage_integration)
1696+
result: List[Tuple[String, logic.Value]] = list[Tuple[String, logic.Value]]()
1697+
if si.provider != "":
1698+
builtin.list_push(result, builtin.tuple("provider", _make_value_string(si.provider)))
1699+
if si.azure_sas_token != "":
1700+
builtin.list_push(result, builtin.tuple("azure_sas_token", _make_value_string("***")))
1701+
if si.s3_region != "":
1702+
builtin.list_push(result, builtin.tuple("s3_region", _make_value_string(si.s3_region)))
1703+
if si.s3_access_key_id != "":
1704+
builtin.list_push(result, builtin.tuple("s3_access_key_id", _make_value_string("***")))
1705+
if si.s3_secret_access_key != "":
1706+
builtin.list_push(result, builtin.tuple("s3_secret_access_key", _make_value_string("***")))
1707+
return builtin.some(builtin.list_sort(result))
1708+
1709+
16611710

16621711
def deconstruct_betree_info_config(msg: logic.BeTreeInfo) -> List[Tuple[String, logic.Value]]:
16631712
result: List[Tuple[String, logic.Value]] = list[Tuple[String, logic.Value]]()

proto/relationalai/lqp/v1/logic.proto

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -276,6 +276,18 @@ message BeTreeLocator {
276276
int64 tree_height = 3;
277277
}
278278

279+
message StorageIntegration {
280+
string provider = 1; // "azure" or "s3"
281+
282+
// Options for azure
283+
string azure_sas_token = 2;
284+
285+
// Options for s3
286+
string s3_region = 3;
287+
string s3_access_key_id = 4;
288+
string s3_secret_access_key = 5;
289+
}
290+
279291
message CSVData {
280292
CSVLocator locator = 1;
281293
CSVConfig config = 2;
@@ -314,6 +326,9 @@ message CSVConfig {
314326

315327
// Partitioning (for export)
316328
int64 partition_size_mb = 12;
329+
330+
// Storage integration (credentials for private buckets)
331+
optional StorageIntegration storage_integration = 13;
317332
}
318333

319334
message IcebergData {

0 commit comments

Comments
 (0)