From ee8e0cf7982b8e72a69b43c2414859fd90b02805 Mon Sep 17 00:00:00 2001 From: TivonB-AI2 <124182151+TivonB-AI2@users.noreply.github.com> Date: Mon, 10 Aug 2026 18:31:25 -0400 Subject: [PATCH] Resolve conflict in cherry-pick of 5053116c95fb7f8ba45f99efaa8a87cf356fb010 and change the commit message --- integrations/Gemfile.lock | 217 ++- .../lib/multiwoven/integrations/rollout.rb | 4 + .../integrations/source/one_drive/client.rb | 640 ++++++++ .../multiwoven/integrations/rollout_spec.rb | 32 + .../source/one_drive/client_spec.rb | 1305 +++++++++++++++++ 5 files changed, 2128 insertions(+), 70 deletions(-) create mode 100644 integrations/lib/multiwoven/integrations/source/one_drive/client.rb create mode 100644 integrations/spec/multiwoven/integrations/source/one_drive/client_spec.rb diff --git a/integrations/Gemfile.lock b/integrations/Gemfile.lock index 1a68732e7..c9651c390 100644 --- a/integrations/Gemfile.lock +++ b/integrations/Gemfile.lock @@ -7,7 +7,11 @@ GIT PATH remote: . specs: +<<<<<<< HEAD multiwoven-integrations (0.35.2) +======= + multiwoven-integrations (0.39.1) +>>>>>>> 5053116c9 (chore(CE): Allow Table Selector for One Drive when recursive (#2146)) MailchimpMarketing activesupport async-websocket @@ -50,7 +54,7 @@ GEM MailchimpMarketing (3.0.80) excon (>= 0.76.0, < 1) json (~> 2.1, >= 2.1.0) - activesupport (8.1.3) + activesupport (8.1.3.1) base64 bigdecimal concurrent-ruby (~> 1.0, >= 1.3.1) @@ -66,6 +70,7 @@ GEM addressable (2.8.9) public_suffix (>= 2.0.2, < 8.0) afm (1.0.0) +<<<<<<< HEAD ast (2.4.2) async (2.11.0) console (~> 1.25, >= 1.25.2) @@ -79,6 +84,11 @@ GEM websocket-driver (~> 0.7.0) aws-eventstream (1.3.0) aws-partitions (1.931.0) +======= + ast (2.4.3) + aws-eventstream (1.4.0) + aws-partitions (1.1240.0) +>>>>>>> 5053116c9 (chore(CE): Allow Table Selector for One Drive when recursive (#2146)) aws-sdk-athena (1.83.0) aws-sdk-core (~> 3, >= 3.193.0) aws-sigv4 (~> 1.1) @@ -118,7 +128,11 @@ GEM bigdecimal (3.3.1) builder (3.3.0) byebug (11.1.3) +<<<<<<< HEAD concurrent-ruby (1.3.6) +======= + concurrent-ruby (1.3.8) +>>>>>>> 5053116c9 (chore(CE): Allow Table Selector for One Drive when recursive (#2146)) connection_pool (3.0.2) console (1.25.2) fiber-annotation @@ -172,8 +186,9 @@ GEM zeitwerk (~> 2.6) duckdb (0.10.3.0) bigdecimal (>= 3.1.4) - ethon (0.16.0) + ethon (0.18.0) ffi (>= 1.15.0) + logger excon (0.112.0) faraday (2.8.1) base64 @@ -181,9 +196,10 @@ GEM ruby2_keywords (>= 0.0.4) faraday-follow_redirects (0.3.0) faraday (>= 1, < 3) - faraday-mashify (0.1.1) + faraday-mashify (1.0.2) faraday (~> 2.0) hashie +<<<<<<< HEAD faraday-multipart (1.0.4) multipart-post (~> 2) faraday-net_http (3.0.2) @@ -194,76 +210,105 @@ GEM fiber-annotation (0.2.0) fiber-local (1.1.0) fiber-storage +======= + faraday-multipart (1.2.0) + multipart-post (~> 2.0) + faraday-net_http (3.4.4) + net-http (~> 0.5) + faraday-retry (2.4.0) + faraday (~> 2.0) + ffi (1.17.4-arm64-darwin) + ffi (1.17.4-x64-mingw-ucrt) + ffi (1.17.4-x86_64-darwin) + ffi (1.17.4-x86_64-linux-gnu) +>>>>>>> 5053116c9 (chore(CE): Allow Table Selector for One Drive when recursive (#2146)) fiber-storage (0.1.1) - gapic-common (0.22.0) + gapic-common (1.3.0) faraday (>= 1.9, < 3.a) faraday-retry (>= 1.0, < 3.a) - google-protobuf (>= 3.25, < 5.a) + google-cloud-env (~> 2.2) + google-logging-utils (~> 0.1) + google-protobuf (~> 4.26) googleapis-common-protos (~> 1.6) googleapis-common-protos-types (~> 1.15) - googleauth (~> 1.11) - grpc (~> 1.65) - git (4.3.2) + googleauth (~> 1.12) + grpc (~> 1.66) + git (5.0.5) activesupport (>= 5.0) addressable (~> 2.8) process_executer (~> 4.0) rchardet (~> 1.9) - gli (2.21.1) - google-apis-bigquery_v2 (0.70.0) + gli (2.22.2) + ostruct + google-apis-bigquery_v2 (0.106.0) google-apis-core (>= 0.15.0, < 2.a) - google-apis-core (0.15.0) - addressable (~> 2.5, >= 2.5.1) - googleauth (~> 1.9) - httpclient (>= 2.8.1, < 3.a) - mini_mime (~> 1.0) + google-apis-core (1.2.5) + addressable (~> 2.9) + faraday (~> 2.13) + faraday-follow_redirects (~> 0.3) + googleauth (~> 1.14) + mini_mime (~> 1.1) + multi_json (~> 1.11) representable (~> 3.0) - retriable (>= 2.0, < 4.a) - rexml + retriable (>= 3.1, < 5.0) google-apis-drive_v3 (0.67.0) google-apis-core (>= 0.15.0, < 2.a) - google-apis-sheets_v4 (0.31.0) - google-apis-core (>= 0.14.0, < 2.a) - google-cloud-ai_platform-v1 (0.50.0) - gapic-common (>= 0.21.1, < 2.a) + google-apis-sheets_v4 (0.48.0) + google-apis-core (>= 0.15.0, < 2.a) + google-cloud-ai_platform-v1 (1.47.0) + gapic-common (~> 1.3) google-cloud-errors (~> 1.0) - google-cloud-location (>= 0.7, < 2.a) - google-iam-v1 (>= 0.7, < 2.a) - google-cloud-bigquery (1.49.0) + google-cloud-location (~> 1.0) + google-iam-v1 (~> 1.3) + google-cloud-bigquery (1.64.0) + bigdecimal (>= 3.0, < 5) concurrent-ruby (~> 1.0) - google-apis-bigquery_v2 (~> 0.62) - google-apis-core (~> 0.13) + google-apis-bigquery_v2 (~> 0.71) + google-apis-core (>= 0.18, < 2) google-cloud-core (~> 1.6) googleauth (~> 1.9) mini_mime (~> 1.0) - google-cloud-core (1.7.0) + google-cloud-core (1.9.0) google-cloud-env (>= 1.0, < 3.a) google-cloud-errors (~> 1.0) +<<<<<<< HEAD google-cloud-env (2.1.1) +======= + google-cloud-env (2.4.0) + base64 (~> 0.2) +>>>>>>> 5053116c9 (chore(CE): Allow Table Selector for One Drive when recursive (#2146)) faraday (>= 1.0, < 3.a) - google-cloud-errors (1.4.0) - google-cloud-location (0.8.1) - gapic-common (>= 0.21.1, < 2.a) + google-cloud-errors (1.7.0) + google-cloud-location (1.5.1) + gapic-common (~> 1.3) google-cloud-errors (~> 1.0) - google-iam-v1 (1.0.1) - gapic-common (>= 0.21.1, < 2.a) + google-iam-v1 (1.7.1) + gapic-common (~> 1.3) google-cloud-errors (~> 1.0) +<<<<<<< HEAD grpc-google-iam-v1 (~> 1.1) google-protobuf (4.30.2-arm64-darwin) +======= + grpc-google-iam-v1 (~> 1.11) + google-logging-utils (0.2.0) + google-protobuf (4.35.1-arm64-darwin) +>>>>>>> 5053116c9 (chore(CE): Allow Table Selector for One Drive when recursive (#2146)) bigdecimal - rake (>= 13) - google-protobuf (4.30.2-x64-mingw-ucrt) + rake (~> 13.3) + google-protobuf (4.35.1-x64-mingw-ucrt) bigdecimal - rake (>= 13) - google-protobuf (4.30.2-x86_64-darwin) + rake (~> 13.3) + google-protobuf (4.35.1-x86_64-darwin) bigdecimal - rake (>= 13) - google-protobuf (4.30.2-x86_64-linux) + rake (~> 13.3) + google-protobuf (4.35.1-x86_64-linux-gnu) bigdecimal - rake (>= 13) - googleapis-common-protos (1.6.0) - google-protobuf (>= 3.18, < 5.a) - googleapis-common-protos-types (~> 1.7) + rake (~> 13.3) + googleapis-common-protos (1.9.0) + google-protobuf (~> 4.26) + googleapis-common-protos-types (~> 1.21) grpc (~> 1.41) +<<<<<<< HEAD googleapis-common-protos-types (1.15.0) google-protobuf (>= 3.18, < 5.a) googleauth (1.11.0) @@ -271,7 +316,17 @@ GEM google-cloud-env (~> 2.1) jwt (>= 1.4, < 3.0) multi_json (~> 1.11) +======= + googleapis-common-protos-types (1.23.0) + google-protobuf (~> 4.26) + googleauth (1.17.3) + faraday (>= 1.0, < 3.a) + google-cloud-env (~> 2.2) + google-logging-utils (~> 0.1) + jwt (>= 1.4, < 4.0) +>>>>>>> 5053116c9 (chore(CE): Allow Table Selector for One Drive when recursive (#2146)) os (>= 0.9, < 2.0) + pstore (~> 0.1) signet (>= 0.16, < 2.a) graphlient (0.8.0) faraday (~> 2.0) @@ -283,21 +338,21 @@ GEM graphql-client (0.26.0) activesupport (>= 3.0) graphql (>= 1.13.0) - grpc (1.66.0-arm64-darwin) + grpc (1.83.0-arm64-darwin) google-protobuf (>= 3.25, < 5.0) googleapis-common-protos-types (~> 1.0) - grpc (1.66.0-x64-mingw-ucrt) + grpc (1.83.0-x64-mingw-ucrt) google-protobuf (>= 3.25, < 5.0) googleapis-common-protos-types (~> 1.0) - grpc (1.66.0-x86_64-darwin) + grpc (1.83.0-x86_64-darwin) google-protobuf (>= 3.25, < 5.0) googleapis-common-protos-types (~> 1.0) - grpc (1.66.0-x86_64-linux) + grpc (1.83.0-x86_64-linux-gnu) google-protobuf (>= 3.25, < 5.0) googleapis-common-protos-types (~> 1.0) - grpc-google-iam-v1 (1.8.0) + grpc-google-iam-v1 (1.12.0) google-protobuf (>= 3.18, < 5.a) - googleapis-common-protos (~> 1.4) + googleapis-common-protos (~> 1.9.0) grpc (~> 1.41) hashdiff (1.1.0) hashery (2.1.2) @@ -306,11 +361,10 @@ GEM csv mini_mime (>= 1.0.0) multi_xml (>= 0.5.2) - httpclient (2.8.3) hubspot-api-client (17.2.0) json (~> 2.1, >= 2.1.0) typhoeus (~> 1.4.0) - i18n (1.14.8) + i18n (1.15.2) concurrent-ruby (~> 1.0) ice_nine (0.11.2) inflection (1.0.0) @@ -323,15 +377,20 @@ GEM multi_json (~> 1.15.0) sorbet-runtime jmespath (1.6.2) +<<<<<<< HEAD json (2.7.2) +======= + json (2.21.2) +>>>>>>> 5053116c9 (chore(CE): Allow Table Selector for One Drive when recursive (#2146)) jsonpath (1.1.5) multi_json jwt (2.8.1) base64 - language_server-protocol (3.17.0.3) + language_server-protocol (3.17.0.6) + lint_roller (1.1.0) logger (1.7.0) mini_mime (1.1.5) - minitest (6.0.3) + minitest (6.0.6) drb (~> 2.0) prism (~> 1.5) multi_json (1.15.0) @@ -351,8 +410,9 @@ GEM nokogiri (1.18.9-x86_64-linux-gnu) racc (~> 1.4) os (1.1.4) - parallel (1.24.0) - parser (3.3.1.0) + ostruct (0.6.3) + parallel (1.28.0) + parser (3.3.12.0) ast (~> 2.4.1) racc pdf-reader (2.15.0) @@ -368,14 +428,15 @@ GEM dry-validation (~> 1.10) httparty (>= 0.22.0) prism (1.9.0) - process_executer (4.0.2) + process_executer (4.0.4) track_open_instances (~> 0.1) + pstore (0.2.1) public_suffix (7.0.5) racc (1.8.1) rainbow (3.1.1) - rake (13.2.1) - rchardet (1.10.0) - regexp_parser (2.9.2) + rake (13.4.2) + rchardet (1.10.2) + regexp_parser (2.12.0) representable (3.2.0) declarative (< 0.1.0) trailblazer-option (>= 0.1.1, < 0.2.0) @@ -387,8 +448,13 @@ GEM faraday-net_http (< 4.0.0) hashie (>= 1.2.0, < 6.0) jwt (>= 1.5.6) +<<<<<<< HEAD retriable (3.1.2) rexml (3.4.0) +======= + retriable (4.2.0) + rexml (3.4.4) +>>>>>>> 5053116c9 (chore(CE): Allow Table Selector for One Drive when recursive (#2146)) rsa-pem-from-mod-exp (0.1.0) rspec (3.13.0) rspec-core (~> 3.13.0) @@ -403,19 +469,20 @@ GEM diff-lcs (>= 1.2.0, < 2.0) rspec-support (~> 3.13.0) rspec-support (3.13.1) - rubocop (1.63.5) + rubocop (1.89.0) json (~> 2.3) - language_server-protocol (>= 3.17.0) - parallel (~> 1.10) + language_server-protocol (~> 3.17.0.2) + lint_roller (~> 1.1.0) + parallel (>= 1.10) parser (>= 3.3.0.2) rainbow (>= 2.2.2, < 4.0) - regexp_parser (>= 1.8, < 3.0) - rexml (>= 3.2.5, < 4.0) - rubocop-ast (>= 1.31.1, < 2.0) + regexp_parser (>= 2.9.3, < 3.0) + rubocop-ast (>= 1.49.0, < 2.0) ruby-progressbar (~> 1.7) - unicode-display_width (>= 2.4.0, < 3.0) - rubocop-ast (1.31.3) - parser (>= 3.3.1.0) + unicode-display_width (>= 2.4.0, < 4.0) + rubocop-ast (1.50.0) + parser (>= 3.3.7.2) + prism (~> 1.7) ruby-limiter (2.3.0) ruby-oci8 (2.2.12) ruby-oci8 (2.2.12-x64-mingw-ucrt) @@ -426,23 +493,31 @@ GEM securerandom (0.4.1) sequel (5.80.0) bigdecimal +<<<<<<< HEAD signet (0.19.0) addressable (~> 2.8) faraday (>= 0.17.5, < 3.a) jwt (>= 1.5, < 3.0) multi_json (~> 1.10) +======= + signet (0.22.0) + addressable (~> 2.8) + faraday (>= 0.17.5, < 3.a) + jwt (>= 1.5, < 4.0) +>>>>>>> 5053116c9 (chore(CE): Allow Table Selector for One Drive when recursive (#2146)) simplecov (0.22.0) docile (~> 1.1) simplecov-html (~> 0.11) simplecov_json_formatter (~> 0.1) simplecov-html (0.12.3) simplecov_json_formatter (0.1.4) - slack-ruby-client (2.3.0) - faraday (>= 2.0) + slack-ruby-client (3.2.0) + faraday (>= 2.0.1) faraday-mashify faraday-multipart gli hashie + logger sorbet-runtime (0.5.11414) stripe (11.4.0) timers (4.3.5) @@ -457,7 +532,9 @@ GEM tzinfo (2.0.6) concurrent-ruby (~> 1.0) uber (0.1.0) - unicode-display_width (2.5.0) + unicode-display_width (3.2.0) + unicode-emoji (~> 4.1) + unicode-emoji (4.2.0) uri (1.1.1) weaviate-ruby (0.9.2) faraday (>= 2.0.1, < 3.0) diff --git a/integrations/lib/multiwoven/integrations/rollout.rb b/integrations/lib/multiwoven/integrations/rollout.rb index 51ad50838..192a18a4b 100644 --- a/integrations/lib/multiwoven/integrations/rollout.rb +++ b/integrations/lib/multiwoven/integrations/rollout.rb @@ -2,7 +2,11 @@ module Multiwoven module Integrations +<<<<<<< HEAD VERSION = "0.35.2" +======= + VERSION = "0.39.1" +>>>>>>> 5053116c9 (chore(CE): Allow Table Selector for One Drive when recursive (#2146)) ENABLED_SOURCES = %w[ Snowflake diff --git a/integrations/lib/multiwoven/integrations/source/one_drive/client.rb b/integrations/lib/multiwoven/integrations/source/one_drive/client.rb new file mode 100644 index 000000000..8556f6f50 --- /dev/null +++ b/integrations/lib/multiwoven/integrations/source/one_drive/client.rb @@ -0,0 +1,640 @@ +# frozen_string_literal: true + +module Multiwoven::Integrations::Source + module OneDrive + include Multiwoven::Integrations::Core + class Client < UnstructuredSourceConnector + SPREADSHEET_EXTENSIONS = %w[.csv .xlsx .xls .xlsm].freeze + EXPIRED_ACCESS_TOKEN_ERROR_CODE = "InvalidAuthenticationToken" + + def check_connection(connection_config) + connection_config = connection_config.with_indifferent_access + if unstructured_data?(connection_config) + create_connection(connection_config) + fetch_list_items + else + conn = create_connection(connection_config) + @sync_id = "check_connection" + files = spreadsheet_files(fetch_list_items) + raise StandardError, "No spreadsheet files found" if files.empty? + + files.each { |file| describe_spreadsheet_file(conn, file) } + end + + success_status + rescue StandardError, NotImplementedError => e + handle_exception(e, { + context: "ONE_DRIVE:CHECK_CONNECTION:EXCEPTION", + type: "error" + }) + failure_status(e) + end + + def discover(connection_config) + connection_config = connection_config.with_indifferent_access + + streams = if unstructured_data?(connection_config) + [create_unstructured_stream] + else + conn = create_connection(connection_config) + @sync_id = "discover" + files = spreadsheet_files(fetch_list_items) + raise StandardError, "No spreadsheet files found" if files.empty? + + files.map { |file| discover_stream_for_file(conn, file) } + end + catalog = Catalog.new(streams: streams) + catalog.to_multiwoven_message + rescue StandardError => e + handle_exception(e, { context: "ONE_DRIVE:DISCOVER:EXCEPTION", type: "error" }) + end + + def read(sync_config) + connection_config = sync_config.source.connection_specification.with_indifferent_access + @connector_instance = sync_config&.source&.connector_instance + + return handle_unstructured_data(sync_config) if unstructured_data?(connection_config) + + conn = create_connection(connection_config) + + @connection_config = connection_config + @sync_id = sync_config.sync_id + + query = sync_config.model.query + query = batched_query(query, sync_config.limit, sync_config.offset) unless sync_config.limit.nil? && sync_config.offset.nil? + query(conn, query) + rescue StandardError => e + handle_exception(e, { + context: "ONE_DRIVE:READ:EXCEPTION", + type: "error", + sync_id: sync_config.sync_id, + sync_run_id: sync_config.sync_run_id + }) + end + + private + + def load_connection_config(connection_config) + @user_name = connection_config[:user_name] + @tenant_id = connection_config[:tenant_id] + @client_id = connection_config[:client_id] + @client_secret = connection_config[:client_secret] + @data_type = connection_config[:data_type] + @file_name = connection_config[:file_name] + @share_url = connection_config[:share_url] + @is_recursive = [true, "true"].include?(connection_config[:is_recursive]) + stored_token = @connector_instance&.configuration&.dig("access_token") + @access_token = stored_token.presence || refresh_access_token + end + + def create_connection(connection_config) + raise ArgumentError, "User Name or Share URL is required" unless connection_config[:share_url].present? || connection_config[:user_name].present? + + load_connection_config(connection_config) + @sync_id ||= "preview" + + if @share_url.present? + @drive_id = shared_folder_reference[:drive_id] + else + response = microsoft_graph_request(user_drive_url) + raise graph_api_error(response.body) unless success?(response) + + @drive_id = JSON.parse(response.body)["id"] + end + + return if @data_type.to_s == "unstructured" + + duckdb_connection + end + + def refresh_access_token + @access_token = fetch_access_token + persist_access_token(@access_token) + @access_token + end + + def persist_access_token(token) + return unless @connector_instance&.configuration + + config = @connector_instance.configuration + config = {} unless config.is_a?(Hash) + @connector_instance.update!(configuration: config.merge("access_token" => token)) + end + + def microsoft_graph_request(url) + response = graph_http_get(url) + return response unless expired_access_token_error?(response.body) + + refresh_access_token + graph_http_get(url) + end + + def graph_http_get(url) + Multiwoven::Integrations::Core::HttpClient.request( + url, + HTTP_GET, + headers: auth_headers(@access_token) + ) + end + + def fetch_access_token + response = Multiwoven::Integrations::Core::HttpClient.request( + format(MICROSOFT_GRAPH_TOKEN_URL, tenant_id: @tenant_id), + HTTP_POST, + payload: form_urlencoded_payload( + client_id: @client_id, + client_secret: @client_secret, + scope: MICROSOFT_GRAPH_SCOPE, + grant_type: "client_credentials" + ), + headers: { + "Content-Type" => "application/x-www-form-urlencoded" + } + ) + raise graph_api_error(response.body) unless success?(response) + + JSON.parse(response.body)["access_token"] + end + + def handle_unstructured_data(sync_config) + connection_config = sync_config.source.connection_specification.with_indifferent_access + command = sync_config.model.query.strip + create_connection(connection_config) + + case command + when LIST_FILES_CMD + list_files_in_folder(connection_config) + when /^#{DOWNLOAD_FILE_CMD}\s+(.+)$/ + file_name = ::Regexp.last_match(1).strip + file_name = file_name.gsub(/^["']|["']$/, "") + download_unstructured_file(connection_config, file_name, sync_config.sync_id) + else + raise ArgumentError, "Invalid command. Supported commands: #{LIST_FILES_CMD}, #{DOWNLOAD_FILE_CMD} " + end + end + + def list_files_in_folder(_connection_config) + files_in_folder.map do |file| + relative_path = relative_file_path(file) + RecordMessage.new( + data: { + element_id: file["id"], + file_name: file["name"], + file_path: relative_path, + size: file["size"], + file_type: File.extname(file["name"]).sub(".", ""), + created_date: file["createdDateTime"], + modified_date: file["lastModifiedDateTime"], + text: "" + }, + emitted_at: Time.now.to_i + ).to_multiwoven_message + end + end + + def download_unstructured_file(_connection_config, file_path, sync_id) + lookup_path = resolve_download_file_name(file_path) + file_item = find_file_item(lookup_path) + raise StandardError, "File not found." if file_item.nil? + + file_name = file_item["name"] + relative_path = relative_file_path(file_item) + local_path = download_file_to_local( + file_name, + sync_id, + item_id: file_item["id"], + drive_id: file_item.dig("parentReference", "driveId") + ) + + [RecordMessage.new( + data: { + element_id: file_item["id"], + local_path: local_path, + file_name: file_name, + file_path: relative_path, + size: file_item["size"], + file_type: File.extname(file_name).sub(".", ""), + created_date: file_item["createdDateTime"], + modified_date: file_item["lastModifiedDateTime"], + text: "" + }, + emitted_at: Time.now.to_i + ).to_multiwoven_message] + end + + def files_in_folder + records = fetch_list_items + records["value"].select do |item| + item["folder"].blank? && matching_file_name?(item["name"], relative_file_path(item)) + end + end + + def find_file_item(lookup_path) + files = files_in_folder + exact_match = files.find do |item| + relative_path = relative_file_path(item) + relative_path == lookup_path + end + return exact_match if exact_match + return if lookup_path.include?("/") + + files.find { |item| item["name"] == lookup_path } + end + + def relative_file_path(file) + file["relative_path"].presence || file["name"] + end + + def resolve_download_file_name(file_path) + return file_path.to_s.strip unless file_path.to_s.start_with?("http") + + @file_name.to_s.strip.presence || File.basename(file_path) + end + + def matching_file_name?(name, relative_path = nil) + configured_name = @file_name.to_s.strip + configured_name.blank? || configured_name == name || configured_name == relative_path + end + + def discover_stream_for_file(conn, file) + describe_results = describe_spreadsheet_file(conn, file) + columns = build_discover_columns(describe_results) + + Multiwoven::Integrations::Protocol::Stream.new( + name: stream_name_for(file), + action: StreamAction["fetch"], + json_schema: convert_to_json_schema(columns) + ) + end + + def describe_spreadsheet_file(conn, file) + local_file = nil + file_name = file["name"] + local_file = download_file_to_local( + file_name, + @sync_id, + item_id: file["id"], + drive_id: file.dig("parentReference", "driveId") + ) + duckdb_file = read_local_file(conn, file_name, local_file) + get_results(conn, "DESCRIBE SELECT * FROM #{duckdb_file};") + ensure + cleanup_ephemeral_download(local_file) + end + + def build_discover_columns(describe_results) + describe_results.map do |row| + { + column_name: row["column_name"], + type: column_schema_helper(row["column_type"]) + } + end + end + + # Maps DuckDB column types before convert_to_json_schema. Note that map_type_to_json_schema + # only recognizes "NUMBER" and "vector", so integer/number/boolean here still become "string" + # in the emitted json_schema (inherited from amazon_s3; not a typed-schema connector). + def column_schema_helper(column_type) + case column_type + when "VARCHAR", "BIT", "DATE", "TIME", "TIMESTAMP", "UUID" + "string" + when "DOUBLE" + "number" + when "BIGINT", "HUGEINT", "INTEGER", "SMALLINT" + "integer" + when "BOOLEAN" + "boolean" + end + end + + def query(connection, query) + local_file = nil + file_name = extract_file_name_from_query(query) + local_basename = File.basename(file_name) + local_file = download_spreadsheet_for_query(file_name) + + file = read_local_file(connection, local_basename, local_file) + query = apply_local_file_to_query(query, file) + get_results(connection, query).map do |row| + RecordMessage.new(data: row, emitted_at: Time.now.to_i).to_multiwoven_message + end + ensure + cleanup_ephemeral_download(local_file) + end + + # Prefer Graph item IDs (same as discover) so nested / shared-folder files + # resolve during model preview. Resolve id with one path metadata GET — + # never a full folder listing / recursive BFS. Fall back to path-based download. + def download_spreadsheet_for_query(file_name) + file_item = fetch_drive_item_by_path(file_name) + if file_item&.[]("id") + download_file_to_local( + file_name, + @sync_id, + item_id: file_item["id"], + drive_id: file_item.dig("parentReference", "driveId") + ) + else + download_file_to_local(file_name, @sync_id) + end + end + + def fetch_drive_item_by_path(file_name) + return fetch_shared_item_metadata if @share_url.present? && shared_folder_reference[:is_file] + + response = microsoft_graph_request(single_file_item_url(file_name)) + return unless success?(response) + + JSON.parse(response.body) + rescue StandardError + nil + end + + def extract_file_name_from_query(sql_query) + match = sql_query.match( + /\bFROM\s+(?:[`"]([^`"]+)[`"]|'([^']+)'|([^\s;]+))/i + ) + + match&.captures&.compact&.first || + raise(ArgumentError, "Could not extract file name from query") + end + + def duckdb_connection + conn = DuckDB::Database.open.connect + conn.execute(INSTALL_HTTPFS_QUERY) + conn + end + + def read_local_file(conn, file_name, local_file) + escaped_path = local_file.gsub("'", "''") + + case File.extname(file_name).downcase + when ".csv" + "read_csv_auto('#{escaped_path}')" + when ".xlsx", ".xls", ".xlsm" + conn.execute("INSTALL excel; LOAD excel;") + "read_xlsx('#{escaped_path}')" + else + raise ArgumentError, "Unsupported file type: #{file_name}" + end + end + + def apply_local_file_to_query(sql_query, file) + sql_query.sub(/\bFROM\s+(?:`[^`]+`|"[^"]+"|'[^']+'|[^\s;]+)/i, "FROM #{file}") + end + + def get_results(conn, sql_query) + hash_array_values(conn.query(sql_query)) + end + + def hash_array_values(results) + keys = results.columns.map(&:name) + results.map do |row| + Hash[keys.zip(row)] + end + end + + def download_file_to_local(file_name, sync_id, item_id: nil, drive_id: nil) + local_file = local_download_path(file_name, sync_id) + FileUtils.mkdir_p(File.dirname(local_file)) + + response = fetch_file_content(file_content_url(file_name, item_id: item_id, drive_id: drive_id)) + raise graph_api_error(response.body) unless success?(response) + + File.binwrite(local_file, response.body) + local_file + rescue StandardError => e + raise StandardError, "Failed to download file #{file_name}: #{e.message}" + end + + def local_download_path(file_name, sync_id) + sync_id = sync_id.presence || "preview" + download_path = ENV["FILE_DOWNLOAD_PATH"] + if download_path + File.join(download_path, "syncs", sync_id, File.basename(file_name)) + else + @temp_download_dir ||= Dir.mktmpdir("one_drive_#{sync_id}") + File.join(@temp_download_dir, File.basename(file_name)) + end + end + + def cleanup_ephemeral_download(local_file) + return if local_file.blank? || ENV["FILE_DOWNLOAD_PATH"].present? + return unless ephemeral_download?(local_file) + + File.delete(local_file) if File.exist?(local_file) + end + + def ephemeral_download?(local_file) + @temp_download_dir.present? && local_file.start_with?(@temp_download_dir) + end + + def file_content_url(file_name, item_id: nil, drive_id: nil) + if item_id.present? + resolved_drive_id = drive_id || @drive_id + return "#{drive_item_url(resolved_drive_id, item_id)}/content" + end + + if @share_url.present? && shared_folder_reference[:is_file] + shared = shared_folder_reference + return "#{drive_item_url(shared[:drive_id], shared[:item_id])}/content" + end + + "#{single_file_item_url(file_name)}:/content" + end + + def single_file_item_url(file_name) + # Preserve path separators for nested files; encode other unsafe chars. + encoded_file = file_name.to_s.split("/").map { |segment| URI::DEFAULT_PARSER.escape(segment) }.join("/") + + if @share_url.present? + shared = shared_folder_reference + "#{drive_item_url(shared[:drive_id], shared[:item_id])}:/#{encoded_file}" + else + "#{drive_root_url(@drive_id)}:/#{encoded_file}" + end + end + + def list_items_url + if @share_url.present? + "#{share_item_url}/children" + else + "#{drive_root_url(@drive_id)}/children" + end + end + + def user_drive_url + format(MICROSOFT_GRAPH_USER_DRIVE_URL, user_name: @user_name) + end + + def share_item_url + share_id = encode_sharing_url(@share_url) + format(MICROSOFT_GRAPH_SHARE_ITEM_URL, share_id: share_id) + end + + def drive_item_url(drive_id, item_id) + format(MICROSOFT_GRAPH_DRIVE_ITEM_URL, drive_id: drive_id, item_id: item_id) + end + + def drive_root_url(drive_id) + "#{MICROSOFT_GRAPH_BASE}/drives/#{drive_id}/root" + end + + def fetch_file_content(url) + response = microsoft_graph_request(url) + return response unless response.is_a?(Net::HTTPRedirection) + + Multiwoven::Integrations::Core::HttpClient.request(response["location"], HTTP_GET) + end + + def fetch_list_items + return { "value" => [fetch_single_file_item] } if single_file_mode? + + return { "value" => [fetch_shared_item_metadata] } if @share_url.present? && shared_folder_reference[:is_file] + + collect_files_from_folder(list_items_url, recursive: @is_recursive) + end + + # Lists files under the configured folder. When recursive is true, BFS over + # /children so nested folder files are included in syncs. + def collect_files_from_folder(root_url, recursive: false) + files = [] + queue = [[root_url, ""]] + + until queue.empty? + url, prefix = queue.shift + page = paginated_graph_collection(url) + + page["value"].each do |item| + relative_path = prefix.empty? ? item["name"].to_s : "#{prefix}/#{item["name"]}" + + if item["folder"].present? + next unless recursive + + drive_id = item.dig("parentReference", "driveId") || @drive_id + queue << ["#{drive_item_url(drive_id, item["id"])}/children", relative_path] + else + item["relative_path"] = relative_path + files << item + end + end + end + + { "value" => files } + end + + def fetch_shared_item_metadata + response = microsoft_graph_request(share_item_url) + raise graph_api_error(response.body) unless success?(response) + + JSON.parse(response.body) + end + + def single_file_mode? + @file_name.to_s.strip.present? && @data_type.to_s == "unstructured" + end + + def fetch_single_file_item + return fetch_shared_item_metadata if @share_url.present? && shared_folder_reference[:is_file] + + response = microsoft_graph_request(single_file_item_url(@file_name)) + raise graph_api_error(response.body) unless success?(response) + + item = JSON.parse(response.body) + item["relative_path"] = @file_name.to_s.strip + item + end + + def paginated_graph_collection(url) + items = [] + next_url = url + + loop do + response = microsoft_graph_request(next_url) + raise graph_api_error(response.body) unless success?(response) + + page = JSON.parse(response.body) + items.concat(page["value"] || []) + next_url = page["@odata.nextLink"] + break if next_url.blank? + end + + { "value" => items } + end + + def shared_folder_reference + @shared_folder_reference ||= begin + response = microsoft_graph_request(share_item_url) + raise graph_api_error(response.body) unless success?(response) + + item = JSON.parse(response.body) + drive_id = item.dig("parentReference", "driveId") + item_id = item["id"] + raise StandardError, "Could not resolve shared folder drive reference" if drive_id.blank? || item_id.blank? + + { drive_id: drive_id, item_id: item_id, is_file: item["file"].present? } + end + end + + def spreadsheet_files(records) + records["value"].select do |record| + record["folder"].blank? && + SPREADSHEET_EXTENSIONS.include?(File.extname(record["name"].to_s).downcase) && + matching_file_name?(record["name"], relative_file_path(record)) + end + end + + # Keep the relative path (with extension) in the stream name — TableSelector + # generates `SELECT * FROM ${stream.name}`, and read_local_file keys off + # File.extname. Relative paths also disambiguate nested duplicates. + def stream_name_for(file) + relative_file_path(file) + end + + def encode_sharing_url(url) + encoded = Base64.strict_encode64(url).tr("+/", "-_").delete("=") + "u!#{encoded}" + end + + def graph_api_error(response_body) + parsed = JSON.parse(response_body) + error = parsed["error"] + + message = if error.is_a?(Hash) + "#{error["code"]}: #{error["message"]}" + elsif error.is_a?(String) + description = parsed["error_description"] + description.present? ? "#{error}: #{description}" : error + else + response_body + end + + StandardError.new(message) + rescue JSON::ParserError, TypeError + StandardError.new(response_body.to_s) + end + + def expired_access_token_error?(response_body) + error = JSON.parse(response_body)["error"] + return false unless error.is_a?(Hash) + + error["code"] == EXPIRED_ACCESS_TOKEN_ERROR_CODE + rescue JSON::ParserError + false + end + + # HttpClient.request always calls payload.to_json. + # Microsoft OAuth token endpoints require + # application/x-www-form-urlencoded bodies instead of JSON. + # This wrapper overrides to_json so HttpClient sends a + # form-encoded string rather than a JSON document. + def form_urlencoded_payload(fields) + payload = Object.new + payload.define_singleton_method(:to_json) do |*_args| + URI.encode_www_form(fields) + end + payload + end + end + end +end diff --git a/integrations/spec/multiwoven/integrations/rollout_spec.rb b/integrations/spec/multiwoven/integrations/rollout_spec.rb index ec834064a..de2d8c935 100644 --- a/integrations/spec/multiwoven/integrations/rollout_spec.rb +++ b/integrations/spec/multiwoven/integrations/rollout_spec.rb @@ -1,6 +1,38 @@ # frozen_string_literal: true RSpec.describe Multiwoven::Integrations do +<<<<<<< HEAD +======= + describe "::VERSION" do + it "is a valid semantic version string" do + expect(Multiwoven::Integrations::VERSION).to match(/\A\d+\.\d+\.\d+\z/) + end + + it "matches the currently released version" do + expect(Multiwoven::Integrations::VERSION).to eq("0.39.1") + end + end + + describe "gem packaging" do + let(:gemspec) do + Gem::Specification.load( + File.expand_path("../../../multiwoven-integrations.gemspec", __dir__) + ) + end + + it "keeps the gemspec version in sync with rollout" do + expect(gemspec.version).to eq(Gem::Version.new(Multiwoven::Integrations::VERSION)) + end + + it "depends on duckdb 1.4.3.0 to match libduckdb v1.4.3" do + duckdb = gemspec.runtime_dependencies.find { |dep| dep.name == "duckdb" } + + expect(duckdb).not_to be_nil + expect(duckdb.requirement).to eq(Gem::Requirement.new("= 1.4.3.0")) + end + end + +>>>>>>> 5053116c9 (chore(CE): Allow Table Selector for One Drive when recursive (#2146)) describe "::ENABLED_SOURCES" do let(:enabled_sources) { Multiwoven::Integrations::ENABLED_SOURCES } let(:enabled_destinations) { Multiwoven::Integrations::ENABLED_DESTINATIONS } diff --git a/integrations/spec/multiwoven/integrations/source/one_drive/client_spec.rb b/integrations/spec/multiwoven/integrations/source/one_drive/client_spec.rb new file mode 100644 index 000000000..4a02d4b6d --- /dev/null +++ b/integrations/spec/multiwoven/integrations/source/one_drive/client_spec.rb @@ -0,0 +1,1305 @@ +# frozen_string_literal: true + +RSpec.describe Multiwoven::Integrations::Source::OneDrive::Client do + let(:client) { described_class.new } + let(:duckdb_conn) { instance_double(DuckDB::Connection) } + let(:access_token) { "test-access-token" } + let(:graph_auth_headers) do + { + "Accept" => "application/json", + "Authorization" => "Bearer #{access_token}", + "Content-Type" => "application/json" + } + end + + let(:structured_config) do + { + data_type: "structured", + user_name: "user@example.com", + tenant_id: "tenant-id", + client_id: "client-id", + client_secret: "client-secret" + } + end + + let(:unstructured_config) do + { + data_type: "unstructured", + user_name: "user@example.com", + tenant_id: "tenant-id", + client_id: "client-id", + client_secret: "client-secret", + share_url: "https://example.com/share-link" + } + end + + let(:sync_config) do + { + source: { + name: "OneDrive", + type: "source", + connection_specification: structured_config + }, + destination: { + name: "Sample Destination Connector", + type: "destination", + connection_specification: {} + }, + model: { + name: "sales", + query: "SELECT * FROM sales.csv", + query_type: "raw_sql", + primary_key: "id" + }, + stream: { + name: "sales.csv", + json_schema: {} + }, + sync_mode: "incremental", + destination_sync_mode: "insert", + sync_id: "sync-1", + sync_run_id: nil + } + end + + let(:unstructured_sync_config) do + { + source: { + name: "OneDrive", + type: "source", + connection_specification: unstructured_config + }, + destination: sync_config[:destination], + model: { + name: "files", + query: "list_files", + query_type: "raw_sql", + primary_key: "element_id" + }, + stream: { + name: "unstructured", + json_schema: {} + }, + sync_mode: "incremental", + destination_sync_mode: "insert", + sync_id: "sync-1", + sync_run_id: "run-1" + } + end + + let(:list_items_response) do + { + "value" => [ + { + "id" => "file-1", + "name" => "report.pdf", + "size" => 1024, + "createdDateTime" => "2024-01-01T00:00:00Z", + "lastModifiedDateTime" => "2024-01-02T00:00:00Z" + }, + { + "id" => "folder-1", + "name" => "nested", + "folder" => { "childCount" => 1 } + } + ] + } + end + + let(:spreadsheet_list_response) do + { + "value" => [ + { "id" => "csv-1", "name" => "sales.csv" }, + { "id" => "txt-1", "name" => "notes.txt" }, + { "id" => "xlsx-1", "name" => "report.xlsx" } + ] + } + end + + let(:describe_results) do + [ + { "column_name" => "id", "column_type" => "BIGINT" }, + { "column_name" => "amount", "column_type" => "DOUBLE" } + ] + end + + let(:share_url) { unstructured_config[:share_url] } + let(:share_id) { client.send(:encode_sharing_url, share_url) } + let(:share_item_url) do + format( + Multiwoven::Integrations::Core::Constants::MICROSOFT_GRAPH_SHARE_ITEM_URL, + share_id: share_id + ) + end + let(:user_drive_url) do + format( + Multiwoven::Integrations::Core::Constants::MICROSOFT_GRAPH_USER_DRIVE_URL, + user_name: structured_config[:user_name] + ) + end + + before do + allow(client).to receive(:create_connection) + allow(client).to receive(:fetch_list_items).and_return(list_items_response) + end + + describe "#check_connection" do + context "with structured data" do + before do + allow(client).to receive(:create_connection).and_return(duckdb_conn) + allow(client).to receive(:fetch_list_items).and_return(spreadsheet_list_response) + allow(client).to receive(:describe_spreadsheet_file) + end + + it "returns succeeded after DESCRIBE validates each spreadsheet" do + expect(client).to receive(:describe_spreadsheet_file).with(duckdb_conn, hash_including("name" => "sales.csv")).ordered + expect(client).to receive(:describe_spreadsheet_file).with(duckdb_conn, hash_including("name" => "report.xlsx")).ordered + + result = client.check_connection(structured_config).connection_status + + expect(result.status).to eq("succeeded") + expect(result.message).to be_nil + end + + it "returns failed when no spreadsheet files are found" do + allow(client).to receive(:fetch_list_items).and_return({ "value" => [] }) + + result = client.check_connection(structured_config).connection_status + + expect(result.status).to eq("failed") + expect(result.message).to include("No spreadsheet files found") + end + + it "returns failed when a spreadsheet cannot be read" do + allow(client).to receive(:describe_spreadsheet_file).and_raise(StandardError, "corrupt file") + + result = client.check_connection(structured_config).connection_status + + expect(result.status).to eq("failed") + expect(result.message).to include("corrupt file") + end + end + + context "with unstructured data" do + it "returns succeeded after listing files" do + result = client.check_connection(unstructured_config).connection_status + + expect(result.status).to eq("succeeded") + expect(result.message).to be_nil + end + + it "returns failed when listing files fails" do + allow(client).to receive(:fetch_list_items).and_raise(StandardError, "listing files failed") + + result = client.check_connection(unstructured_config).connection_status + + expect(result.status).to eq("failed") + expect(result.message).to include("listing files failed") + end + end + + context "when create_connection raises" do + it "returns failed with the error message" do + allow(client).to receive(:create_connection).and_raise(StandardError, "Connection failed") + + result = client.check_connection(structured_config).connection_status + + expect(result.status).to eq("failed") + expect(result.message).to include("Connection failed") + end + end + end + + describe "#discover" do + context "with unstructured data" do + it "returns a catalog with a single unstructured stream" do + message = client.discover(unstructured_config) + + expect(message.catalog).to be_a(Multiwoven::Integrations::Protocol::Catalog) + expect(message.catalog.streams).to be_an(Array) + expect(message.catalog.streams.first).to be_a(Multiwoven::Integrations::Protocol::Stream) + expect(message.catalog.streams.first.name).to eq("unstructured") + end + end + + context "with structured data" do + before do + allow(client).to receive(:create_connection).and_return(duckdb_conn) + allow(client).to receive(:fetch_list_items).and_return(spreadsheet_list_response) + allow(client).to receive(:download_file_to_local).and_return("/tmp/file.csv") + allow(client).to receive(:read_local_file).and_return("read_csv_auto('/tmp/file.csv')") + allow(client).to receive(:get_results).and_return(describe_results) + end + + it "returns a stream per spreadsheet using the full file name" do + message = client.discover(structured_config) + stream_names = message.catalog.streams.map(&:name) + + expect(stream_names).to contain_exactly("sales.csv", "report.xlsx") + end + + it "uses relative paths as stream names for nested spreadsheets" do + allow(client).to receive(:fetch_list_items).and_return( + { + "value" => [ + { + "id" => "csv-nested", + "name" => "sales.csv", + "relative_path" => "reports/2024/sales.csv" + } + ] + } + ) + + message = client.discover(structured_config) + + expect(message.catalog.streams.map(&:name)).to eq(["reports/2024/sales.csv"]) + end + + it "builds all-string json_schema from DuckDB DESCRIBE results" do + message = client.discover(structured_config) + schema = message.catalog.streams.first.json_schema + + expect(schema["properties"].keys).to contain_exactly("id", "amount") + expect(schema["properties"]["id"]["type"]).to eq("string") + expect(schema["properties"]["amount"]["type"]).to eq("string") + end + + it "handles exceptions when no spreadsheet files are found" do + allow(client).to receive(:fetch_list_items).and_return({ "value" => [] }) + + expect(client).to receive(:handle_exception).with( + an_instance_of(StandardError), + { context: "ONE_DRIVE:DISCOVER:EXCEPTION", type: "error" } + ) + + client.discover(structured_config) + end + + it "removes ephemeral downloads after each DESCRIBE" do + allow(ENV).to receive(:[]).and_call_original + allow(ENV).to receive(:[]).with("FILE_DOWNLOAD_PATH").and_return(nil) + + temp_dir = Dir.mktmpdir("one_drive_spec") + client.instance_variable_set(:@temp_download_dir, temp_dir) + local_path = File.join(temp_dir, "sales.csv") + allow(client).to receive(:download_file_to_local).and_return(local_path) + File.write(local_path, "col1\n1") + + client.send(:describe_spreadsheet_file, duckdb_conn, spreadsheet_list_response["value"].first) + + expect(File).not_to exist(local_path) + ensure + FileUtils.remove_entry(temp_dir) if temp_dir && Dir.exist?(temp_dir) + end + end + + context "when discovery fails" do + it "handles exceptions" do + allow(unstructured_config).to receive(:with_indifferent_access).and_raise(StandardError, "Discovery failed") + + expect(client).to receive(:handle_exception).with( + an_instance_of(StandardError), + { context: "ONE_DRIVE:DISCOVER:EXCEPTION", type: "error" } + ) + + client.discover(unstructured_config) + end + end + end + + describe "#read" do + context "with structured data" do + before do + allow(client).to receive(:create_connection).and_return(duckdb_conn) + allow(client).to receive(:fetch_drive_item_by_path).and_return(nil) + allow(client).to receive(:download_file_to_local).and_return("/tmp/sales.csv") + allow(client).to receive(:read_local_file).and_return("read_csv_auto('/tmp/sales.csv')") + allow(client).to receive(:get_results).and_return([{ "id" => "1", "amount" => "100" }]) + end + + it "reads records successfully" do + s_config = Multiwoven::Integrations::Protocol::SyncConfig.from_json(sync_config.to_json) + records = client.read(s_config) + + expect(records).to be_an(Array) + expect(records).not_to be_empty + expect(records.first).to be_a(Multiwoven::Integrations::Protocol::MultiwovenMessage) + expect(records.first.record.data).to eq({ "id" => "1", "amount" => "100" }) + end + + it "reads records with batched query" do + s_config = Multiwoven::Integrations::Protocol::SyncConfig.from_json(sync_config.to_json) + s_config.limit = 100 + s_config.offset = 1 + + expect(client).to receive(:get_results).with( + duckdb_conn, + a_string_including("LIMIT 100 OFFSET 1") + ).and_return([{ "id" => "1" }]) + + records = client.read(s_config) + + expect(records).to be_an(Array) + expect(records.first.record.data).to eq({ "id" => "1" }) + end + + it "handles read failures" do + s_config = Multiwoven::Integrations::Protocol::SyncConfig.from_json(sync_config.to_json) + s_config.sync_run_id = "run-2" + allow(client).to receive(:create_connection).and_raise(StandardError, "test error") + + expect(client).to receive(:handle_exception).with( + an_instance_of(StandardError), + { + context: "ONE_DRIVE:READ:EXCEPTION", + type: "error", + sync_id: "sync-1", + sync_run_id: "run-2" + } + ) + + client.read(s_config) + end + + describe "#extract_file_name_from_query" do + it "extracts backtick-quoted file names" do + expect(client.send(:extract_file_name_from_query, "SELECT * FROM `sales.csv`")).to eq("sales.csv") + end + + it "extracts double-quoted file names" do + expect(client.send(:extract_file_name_from_query, 'SELECT * FROM "report.xlsx"')).to eq("report.xlsx") + end + + it "extracts single-quoted file names" do + expect(client.send(:extract_file_name_from_query, "SELECT * FROM 'data.csv'")).to eq("data.csv") + end + + it "extracts bare file names" do + expect(client.send(:extract_file_name_from_query, "SELECT * FROM sales.csv")).to eq("sales.csv") + end + + it "extracts nested relative paths from queries" do + expect(client.send(:extract_file_name_from_query, "SELECT * FROM `reports/2024/sales.csv`")) + .to eq("reports/2024/sales.csv") + end + + it "extracts file names before a terminating semicolon" do + expect(client.send(:extract_file_name_from_query, "SELECT * FROM sales.csv;")).to eq("sales.csv") + end + + it "raises when the file name cannot be extracted" do + expect do + client.send(:extract_file_name_from_query, "SELECT 1") + end.to raise_error(ArgumentError, "Could not extract file name from query") + end + end + + it "uses the basename when reading a nested spreadsheet path" do + sync_config[:model][:query] = "SELECT * FROM `reports/2024/sales.csv`" + s_config = Multiwoven::Integrations::Protocol::SyncConfig.from_json(sync_config.to_json) + allow(client).to receive(:download_file_to_local).and_return("/tmp/sales.csv") + + expect(client).to receive(:read_local_file).with( + duckdb_conn, + "sales.csv", + "/tmp/sales.csv" + ).and_return("read_csv_auto('/tmp/sales.csv')") + + records = client.read(s_config) + + expect(records.first.record.data).to eq({ "id" => "1", "amount" => "100" }) + expect(client).to have_received(:download_file_to_local).with("reports/2024/sales.csv", "sync-1") + end + + it "downloads by Graph item id from a single path metadata lookup" do + allow(client).to receive(:fetch_drive_item_by_path).with("reports/2024/sales.csv").and_return( + { + "id" => "csv-nested", + "name" => "sales.csv", + "parentReference" => { "driveId" => "drive-1" } + } + ) + sync_config[:model][:query] = "SELECT * FROM `reports/2024/sales.csv`" + s_config = Multiwoven::Integrations::Protocol::SyncConfig.from_json(sync_config.to_json) + allow(client).to receive(:download_file_to_local).and_return("/tmp/sales.csv") + + records = client.read(s_config) + + expect(records.first.record.data).to eq({ "id" => "1", "amount" => "100" }) + expect(client).not_to have_received(:fetch_list_items) + expect(client).to have_received(:download_file_to_local).with( + "reports/2024/sales.csv", + "sync-1", + item_id: "csv-nested", + drive_id: "drive-1" + ) + end + + it "downloads by item id for unquoted table-selector queries" do + allow(client).to receive(:fetch_drive_item_by_path).with("sales.csv").and_return( + { + "id" => "csv-1", + "name" => "sales.csv", + "parentReference" => { "driveId" => "drive-1" } + } + ) + s_config = Multiwoven::Integrations::Protocol::SyncConfig.from_json(sync_config.to_json) + allow(client).to receive(:download_file_to_local).and_return("/tmp/sales.csv") + + client.read(s_config) + + expect(client).not_to have_received(:fetch_list_items) + expect(client).to have_received(:download_file_to_local).with( + "sales.csv", + "sync-1", + item_id: "csv-1", + drive_id: "drive-1" + ) + end + + it "falls back to path-based download when path metadata lookup misses" do + allow(client).to receive(:fetch_drive_item_by_path).with("sales.csv").and_return(nil) + s_config = Multiwoven::Integrations::Protocol::SyncConfig.from_json(sync_config.to_json) + allow(client).to receive(:download_file_to_local).and_return("/tmp/sales.csv") + + client.read(s_config) + + expect(client).to have_received(:download_file_to_local).with("sales.csv", "sync-1") + end + end + + context "with unstructured data" do + it "returns records for list_files" do + s_config = Multiwoven::Integrations::Protocol::SyncConfig.from_json(unstructured_sync_config.to_json) + records = client.read(s_config) + + expect(records).to be_an(Array) + expect(records.first).to be_a(Multiwoven::Integrations::Protocol::MultiwovenMessage) + expect(records.first.record.data[:element_id]).to eq("file-1") + expect(records.first.record.data[:file_name]).to eq("report.pdf") + expect(records.first.record.data[:file_path]).to eq("report.pdf") + expect(records.first.record.data[:file_type]).to eq("pdf") + end + + it "returns relative file_path for nested files from list_files" do + allow(client).to receive(:fetch_list_items).and_return( + { + "value" => [ + { + "id" => "nested-1", + "name" => "nested.pdf", + "relative_path" => "docs/nested.pdf", + "size" => 2048, + "createdDateTime" => "2024-01-01T00:00:00Z", + "lastModifiedDateTime" => "2024-01-02T00:00:00Z" + } + ] + } + ) + s_config = Multiwoven::Integrations::Protocol::SyncConfig.from_json(unstructured_sync_config.to_json) + records = client.read(s_config) + + expect(records.first.record.data[:file_name]).to eq("nested.pdf") + expect(records.first.record.data[:file_path]).to eq("docs/nested.pdf") + end + + it "returns records for download_file" do + unstructured_sync_config[:model][:query] = "download_file report.pdf" + s_config = Multiwoven::Integrations::Protocol::SyncConfig.from_json(unstructured_sync_config.to_json) + allow(client).to receive(:download_file_to_local).and_return("/tmp/report.pdf") + + records = client.read(s_config) + + expect(records).to be_an(Array) + expect(records.first.record.data[:local_path]).to eq("/tmp/report.pdf") + expect(records.first.record.data[:file_name]).to eq("report.pdf") + expect(records.first.record.data[:file_path]).to eq("report.pdf") + end + + it "downloads a nested file by relative path" do + allow(client).to receive(:fetch_list_items).and_return( + { + "value" => [ + { + "id" => "nested-1", + "name" => "nested.pdf", + "relative_path" => "docs/nested.pdf", + "size" => 2048, + "createdDateTime" => "2024-01-01T00:00:00Z", + "lastModifiedDateTime" => "2024-01-02T00:00:00Z", + "parentReference" => { "driveId" => "drive-1" } + } + ] + } + ) + unstructured_sync_config[:model][:query] = "download_file docs/nested.pdf" + s_config = Multiwoven::Integrations::Protocol::SyncConfig.from_json(unstructured_sync_config.to_json) + allow(client).to receive(:download_file_to_local).and_return("/tmp/nested.pdf") + + records = client.read(s_config) + + expect(records.first.record.data[:file_name]).to eq("nested.pdf") + expect(records.first.record.data[:file_path]).to eq("docs/nested.pdf") + expect(client).to have_received(:download_file_to_local).with( + "nested.pdf", + "sync-1", + item_id: "nested-1", + drive_id: "drive-1" + ) + end + + it "returns records for download_file with quoted file names" do + unstructured_sync_config[:model][:query] = 'download_file "report.pdf"' + s_config = Multiwoven::Integrations::Protocol::SyncConfig.from_json(unstructured_sync_config.to_json) + allow(client).to receive(:download_file_to_local).and_return("/tmp/report.pdf") + + records = client.read(s_config) + + expect(records.first.record.data[:file_name]).to eq("report.pdf") + end + + it "handles a missing file" do + unstructured_sync_config[:model][:query] = "download_file missing.pdf" + s_config = Multiwoven::Integrations::Protocol::SyncConfig.from_json(unstructured_sync_config.to_json) + s_config.sync_run_id = "run-1" + + expect(client).to receive(:handle_exception).with( + an_instance_of(StandardError), + { + context: "ONE_DRIVE:READ:EXCEPTION", + type: "error", + sync_id: "sync-1", + sync_run_id: "run-1" + } + ) + + client.read(s_config) + end + + it "handles an invalid command" do + unstructured_sync_config[:model][:query] = "invalid_command" + s_config = Multiwoven::Integrations::Protocol::SyncConfig.from_json(unstructured_sync_config.to_json) + s_config.sync_run_id = "run-1" + + expect(client).to receive(:handle_exception).with( + an_instance_of(ArgumentError), + { + context: "ONE_DRIVE:READ:EXCEPTION", + type: "error", + sync_id: "sync-1", + sync_run_id: "run-1" + } + ) + + client.read(s_config) + end + + context "when FILE_DOWNLOAD_PATH is set" do + before do + allow(ENV).to receive(:[]).with("FILE_DOWNLOAD_PATH").and_return("/custom/download/path") + end + + it "returns the configured download path" do + unstructured_sync_config[:model][:query] = "download_file report.pdf" + s_config = Multiwoven::Integrations::Protocol::SyncConfig.from_json(unstructured_sync_config.to_json) + allow(client).to receive(:download_file_to_local).and_return("/custom/download/path/syncs/sync-1/report.pdf") + + records = client.read(s_config) + + expect(records).to be_an(Array) + expect(records.first.record.data[:local_path]).to eq("/custom/download/path/syncs/sync-1/report.pdf") + end + end + end + end + + describe "connection and share URL resolution" do + before do + allow(client).to receive(:create_connection).and_call_original + allow(client).to receive(:refresh_access_token).and_return(access_token) + allow(DuckDB::Database).to receive(:open).and_return( + instance_double(DuckDB::Database, connect: duckdb_conn) + ) + allow(duckdb_conn).to receive(:execute) + end + + describe "#create_connection" do + context "via user drive" do + it "resolves drive_id and returns a DuckDB connection" do + drive_response = instance_double( + Net::HTTPSuccess, + code: "200", + body: { "id" => "drive-from-user" }.to_json + ) + + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + user_drive_url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: graph_auth_headers + ).and_return(drive_response) + + connection = client.send(:create_connection, structured_config) + + expect(connection).to eq(duckdb_conn) + expect(client.instance_variable_get(:@drive_id)).to eq("drive-from-user") + end + + it "defaults sync_id to preview for model preview execute_query path" do + drive_response = instance_double( + Net::HTTPSuccess, + code: "200", + body: { "id" => "drive-from-user" }.to_json + ) + allow(Multiwoven::Integrations::Core::HttpClient).to receive(:request).and_return(drive_response) + + client.send(:create_connection, structured_config) + + expect(client.instance_variable_get(:@sync_id)).to eq("preview") + end + + it "does not overwrite an existing sync_id" do + drive_response = instance_double( + Net::HTTPSuccess, + code: "200", + body: { "id" => "drive-from-user" }.to_json + ) + allow(Multiwoven::Integrations::Core::HttpClient).to receive(:request).and_return(drive_response) + client.instance_variable_set(:@sync_id, "sync-1") + + client.send(:create_connection, structured_config) + + expect(client.instance_variable_get(:@sync_id)).to eq("sync-1") + end + + it "loads is_recursive from the connection config" do + drive_response = instance_double( + Net::HTTPSuccess, + code: "200", + body: { "id" => "drive-from-user" }.to_json + ) + allow(Multiwoven::Integrations::Core::HttpClient).to receive(:request).and_return(drive_response) + + client.send(:create_connection, structured_config.merge(is_recursive: true)) + expect(client.instance_variable_get(:@is_recursive)).to be(true) + + client.send(:create_connection, structured_config.merge(is_recursive: "true")) + expect(client.instance_variable_get(:@is_recursive)).to be(true) + + client.send(:create_connection, structured_config.merge(is_recursive: false)) + expect(client.instance_variable_get(:@is_recursive)).to be(false) + end + end + + context "via share URL" do + let(:share_response) do + instance_double( + Net::HTTPSuccess, + code: "200", + body: { + "id" => "shared-folder-item", + "parentReference" => { "driveId" => "shared-drive-id" }, + "folder" => { "childCount" => 2 } + }.to_json + ) + end + + it "resolves drive_id and skips DuckDB for unstructured data" do + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + share_item_url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: graph_auth_headers + ).and_return(share_response) + + connection = client.send(:create_connection, unstructured_config) + + expect(connection).to be_nil + expect(client.instance_variable_get(:@drive_id)).to eq("shared-drive-id") + end + + it "resolves drive_id when share_url is provided without user_name" do + config_without_user = unstructured_config.dup.tap { |c| c.delete(:user_name) } + + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + share_item_url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: graph_auth_headers + ).and_return(share_response) + + connection = client.send(:create_connection, config_without_user) + + expect(connection).to be_nil + expect(client.instance_variable_get(:@drive_id)).to eq("shared-drive-id") + expect(client.instance_variable_get(:@user_name)).to be_nil + end + end + + context "when neither user_name nor share_url is provided" do + it "raises an error" do + config_without_user_name = structured_config.dup.tap { |c| c.delete(:user_name) } + + expect do + client.send(:create_connection, config_without_user_name) + end.to raise_error(ArgumentError, "User Name or Share URL is required") + end + end + end + + describe "#encode_sharing_url" do + it "encodes a sharing URL into the Microsoft Graph share id format" do + url = "https://contoso.sharepoint.com/:f:/r/sites/Test/Shared%20Documents" + expected = "u!#{Base64.strict_encode64(url).tr("+/", "-_").delete("=")}" + + expect(client.send(:encode_sharing_url, url)).to eq(expected) + end + end + + describe "#shared_folder_reference" do + before do + client.instance_variable_set(:@share_url, share_url) + client.instance_variable_set(:@access_token, access_token) + end + + it "resolves drive and item ids and memoizes the result" do + share_response = instance_double( + Net::HTTPSuccess, + code: "200", + body: { + "id" => "shared-item-id", + "parentReference" => { "driveId" => "shared-drive-id" }, + "file" => { "mimeType" => "application/pdf" } + }.to_json + ) + + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + share_item_url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: graph_auth_headers + ).once.and_return(share_response) + + first_reference = client.send(:shared_folder_reference) + second_reference = client.send(:shared_folder_reference) + + expect(first_reference).to eq( + drive_id: "shared-drive-id", + item_id: "shared-item-id", + is_file: true + ) + expect(second_reference).to eq(first_reference) + end + end + end + + describe "OAuth token handling" do + let(:token_url) do + format( + Multiwoven::Integrations::Core::Constants::MICROSOFT_GRAPH_TOKEN_URL, + tenant_id: structured_config[:tenant_id] + ) + end + let(:expired_token) { "expired-token" } + let(:fresh_token) { "fresh-token" } + let(:expired_response_body) do + { + "error" => { + "code" => "InvalidAuthenticationToken", + "message" => "Lifetime validation failed, the token is expired." + } + }.to_json + end + let(:expired_headers) do + { + "Accept" => "application/json", + "Authorization" => "Bearer #{expired_token}", + "Content-Type" => "application/json" + } + end + let(:fresh_headers) do + { + "Accept" => "application/json", + "Authorization" => "Bearer #{fresh_token}", + "Content-Type" => "application/json" + } + end + + before do + client.instance_variable_set(:@tenant_id, structured_config[:tenant_id]) + client.instance_variable_set(:@client_id, structured_config[:client_id]) + client.instance_variable_set(:@client_secret, structured_config[:client_secret]) + end + + describe "#fetch_access_token" do + it "requests a client-credentials token with form-urlencoded payload" do + token_response = instance_double( + Net::HTTPSuccess, + code: "200", + body: { "access_token" => "new-access-token" }.to_json + ) + expected_payload = URI.encode_www_form( + client_id: structured_config[:client_id], + client_secret: structured_config[:client_secret], + scope: Multiwoven::Integrations::Core::Constants::MICROSOFT_GRAPH_SCOPE, + grant_type: "client_credentials" + ) + + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + token_url, + Multiwoven::Integrations::Core::Constants::HTTP_POST, + payload: satisfy { |payload| payload.to_json == expected_payload }, + headers: { "Content-Type" => "application/x-www-form-urlencoded" } + ).and_return(token_response) + + expect(client.send(:fetch_access_token)).to eq("new-access-token") + end + + it "raises when the token endpoint returns an error" do + error_response = instance_double( + Net::HTTPUnauthorized, + code: "401", + body: { + "error" => "invalid_client", + "error_description" => "Invalid client secret" + }.to_json + ) + + allow(Multiwoven::Integrations::Core::HttpClient).to receive(:request).and_return(error_response) + + expect { client.send(:fetch_access_token) }.to raise_error(StandardError, /invalid_client/) + end + end + + describe "#form_urlencoded_payload" do + it "serializes fields as application/x-www-form-urlencoded via #to_json" do + payload = client.send( + :form_urlencoded_payload, + client_id: "client-id", + client_secret: "client-secret", + scope: "https://graph.microsoft.com/.default", + grant_type: "client_credentials" + ) + + expect(payload.to_json).to eq( + "client_id=client-id&client_secret=client-secret&scope=https%3A%2F%2Fgraph.microsoft.com%2F.default&grant_type=client_credentials" + ) + end + end + + describe "expired token retry" do + let(:url) { "https://graph.microsoft.com/v1.0/test" } + + before do + client.instance_variable_set(:@access_token, expired_token) + end + + it "refreshes the token and retries microsoft_graph_request once" do + expired_response = instance_double(Net::HTTPUnauthorized, code: "401", body: expired_response_body) + success_response = instance_double(Net::HTTPSuccess, code: "200", body: { "value" => [] }.to_json) + + expect(client).to receive(:refresh_access_token).once do + client.instance_variable_set(:@access_token, fresh_token) + fresh_token + end + + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: expired_headers + ).ordered.and_return(expired_response) + + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: fresh_headers + ).ordered.and_return(success_response) + + expect(client.send(:microsoft_graph_request, url).body).to eq(success_response.body) + end + + it "retries paginated listing when the first page returns an expired token" do + allow(client).to receive(:fetch_list_items).and_call_original + client.instance_variable_set(:@drive_id, "drive-1") + client.instance_variable_set(:@share_url, nil) + client.instance_variable_set(:@data_type, "structured") + + first_page_url = "https://graph.microsoft.com/v1.0/drives/drive-1/root/children" + expired_response = instance_double(Net::HTTPUnauthorized, code: "401", body: expired_response_body) + success_response = instance_double( + Net::HTTPSuccess, + code: "200", + body: { "value" => [{ "id" => "csv-1", "name" => "sales.csv" }] }.to_json + ) + + expect(client).to receive(:refresh_access_token).once do + client.instance_variable_set(:@access_token, fresh_token) + fresh_token + end + + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + first_page_url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: expired_headers + ).ordered.and_return(expired_response) + + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + first_page_url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: fresh_headers + ).ordered.and_return(success_response) + + result = client.send(:fetch_list_items) + + expect(result["value"].map { |item| item["name"] }).to eq(["sales.csv"]) + end + end + end + + describe "Graph file operations" do + before do + client.instance_variable_set(:@access_token, access_token) + end + + describe "#fetch_file_content" do + let(:content_url) { "https://graph.microsoft.com/v1.0/drives/drive-1/items/item-1/content" } + let(:redirect_url) { "https://cdn.example.com/download/file" } + + it "follows a 302 redirect to download file content" do + redirect_response = Net::HTTPFound.new("1.1", "302", "Found") + redirect_response["location"] = redirect_url + file_response = instance_double(Net::HTTPSuccess, code: "200", body: "file-binary-content") + + allow(client).to receive(:microsoft_graph_request).with(content_url).and_return(redirect_response) + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + redirect_url, + Multiwoven::Integrations::Core::Constants::HTTP_GET + ).and_return(file_response) + + expect(client.send(:fetch_file_content, content_url).body).to eq("file-binary-content") + end + + it "returns the original response when no redirect is issued" do + success_response = instance_double(Net::HTTPSuccess, code: "200", body: "inline-content") + + allow(client).to receive(:microsoft_graph_request).with(content_url).and_return(success_response) + expect(Multiwoven::Integrations::Core::HttpClient).not_to receive(:request) + + expect(client.send(:fetch_file_content, content_url).body).to eq("inline-content") + end + end + + describe "#fetch_list_items" do + before do + allow(client).to receive(:fetch_list_items).and_call_original + client.instance_variable_set(:@drive_id, "drive-1") + client.instance_variable_set(:@share_url, nil) + client.instance_variable_set(:@data_type, "structured") + client.instance_variable_set(:@is_recursive, false) + end + + it "follows @odata.nextLink and merges all pages" do + first_page_url = "https://graph.microsoft.com/v1.0/drives/drive-1/root/children" + next_page_url = "https://graph.microsoft.com/v1.0/drives/drive-1/root/children?$skiptoken=abc" + + first_response = instance_double( + Net::HTTPSuccess, + code: "200", + body: { + "value" => [{ "id" => "csv-1", "name" => "sales.csv" }], + "@odata.nextLink" => next_page_url + }.to_json + ) + second_response = instance_double( + Net::HTTPSuccess, + code: "200", + body: { + "value" => [{ "id" => "xlsx-1", "name" => "report.xlsx" }] + }.to_json + ) + + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + first_page_url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: graph_auth_headers + ).and_return(first_response) + + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + next_page_url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: graph_auth_headers + ).and_return(second_response) + + result = client.send(:fetch_list_items) + + expect(result["value"].map { |item| item["name"] }).to eq(%w[sales.csv report.xlsx]) + expect(result["value"].map { |item| item["relative_path"] }).to eq(%w[sales.csv report.xlsx]) + end + + it "skips nested folders when is_recursive is false" do + root_url = "https://graph.microsoft.com/v1.0/drives/drive-1/root/children" + root_response = instance_double( + Net::HTTPSuccess, + code: "200", + body: { + "value" => [ + { "id" => "csv-1", "name" => "sales.csv" }, + { + "id" => "folder-1", + "name" => "nested", + "folder" => { "childCount" => 1 }, + "parentReference" => { "driveId" => "drive-1" } + } + ] + }.to_json + ) + + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + root_url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: graph_auth_headers + ).once.and_return(root_response) + + result = client.send(:fetch_list_items) + + expect(result["value"].map { |item| item["name"] }).to eq(["sales.csv"]) + expect(result["value"].first["relative_path"]).to eq("sales.csv") + end + + it "recursively lists nested folder files when is_recursive is true" do + client.instance_variable_set(:@is_recursive, true) + root_url = "https://graph.microsoft.com/v1.0/drives/drive-1/root/children" + nested_url = "https://graph.microsoft.com/v1.0/drives/drive-1/items/folder-1/children" + + root_response = instance_double( + Net::HTTPSuccess, + code: "200", + body: { + "value" => [ + { "id" => "csv-1", "name" => "sales.csv" }, + { + "id" => "folder-1", + "name" => "nested", + "folder" => { "childCount" => 1 }, + "parentReference" => { "driveId" => "drive-1" } + } + ] + }.to_json + ) + nested_response = instance_double( + Net::HTTPSuccess, + code: "200", + body: { + "value" => [ + { "id" => "pdf-1", "name" => "notes.pdf" } + ] + }.to_json + ) + + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + root_url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: graph_auth_headers + ).ordered.and_return(root_response) + + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + nested_url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: graph_auth_headers + ).ordered.and_return(nested_response) + + result = client.send(:fetch_list_items) + + expect(result["value"].map { |item| item["name"] }).to eq(%w[sales.csv notes.pdf]) + expect(result["value"].map { |item| item["relative_path"] }).to eq(%w[sales.csv nested/notes.pdf]) + end + + context "with a configured file_name" do + let(:file_url) { "https://graph.microsoft.com/v1.0/drives/drive-1/root:/report.pdf" } + let(:file_response) do + instance_double( + Net::HTTPSuccess, + code: "200", + body: { + "id" => "file-1", + "name" => "report.pdf", + "size" => 1024, + "createdDateTime" => "2024-01-01T00:00:00Z", + "lastModifiedDateTime" => "2024-01-02T00:00:00Z" + }.to_json + ) + end + + before do + client.instance_variable_set(:@data_type, "unstructured") + client.instance_variable_set(:@file_name, "report.pdf") + end + + it "fetches only the configured file" do + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + file_url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: graph_auth_headers + ).and_return(file_response) + + result = client.send(:fetch_list_items) + + expect(result["value"].map { |item| item["name"] }).to eq(["report.pdf"]) + end + end + + context "with a share URL" do + let(:children_url) { "#{share_item_url}/children" } + + before do + client.instance_variable_set(:@share_url, share_url) + client.instance_variable_set(:@data_type, "unstructured") + end + + it "lists children when the shared link points to a folder" do + folder_metadata_response = instance_double( + Net::HTTPSuccess, + code: "200", + body: { + "id" => "shared-folder-item", + "parentReference" => { "driveId" => "shared-drive-id" }, + "folder" => { "childCount" => 1 } + }.to_json + ) + children_response = instance_double( + Net::HTTPSuccess, + code: "200", + body: { "value" => [{ "id" => "file-1", "name" => "report.pdf" }] }.to_json + ) + + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + share_item_url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: graph_auth_headers + ).ordered.and_return(folder_metadata_response) + + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + children_url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: graph_auth_headers + ).ordered.and_return(children_response) + + result = client.send(:fetch_list_items) + + expect(result["value"].map { |item| item["name"] }).to eq(["report.pdf"]) + end + + it "returns shared file metadata when the shared link points to a file" do + file_metadata_response = instance_double( + Net::HTTPSuccess, + code: "200", + body: { + "id" => "shared-file-item", + "name" => "shared-report.pdf", + "parentReference" => { "driveId" => "shared-drive-id" }, + "file" => { "mimeType" => "application/pdf" } + }.to_json + ) + + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + share_item_url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: graph_auth_headers + ).twice.and_return(file_metadata_response) + + result = client.send(:fetch_list_items) + + expect(result["value"].map { |item| item["name"] }).to eq(["shared-report.pdf"]) + end + end + end + + describe "#single_file_item_url" do + before do + client.instance_variable_set(:@drive_id, "drive-1") + client.instance_variable_set(:@share_url, nil) + end + + it "preserves path separators for nested file paths" do + expect(client.send(:single_file_item_url, "reports/2024/sales.csv")) + .to eq("https://graph.microsoft.com/v1.0/drives/drive-1/root:/reports/2024/sales.csv") + end + + it "encodes unsafe characters in path segments" do + expect(client.send(:single_file_item_url, "my reports/q1 sales.csv")) + .to eq("https://graph.microsoft.com/v1.0/drives/drive-1/root:/my%20reports/q1%20sales.csv") + end + end + + describe "#fetch_drive_item_by_path" do + before do + client.instance_variable_set(:@drive_id, "drive-1") + client.instance_variable_set(:@share_url, nil) + client.instance_variable_set(:@access_token, access_token) + end + + it "loads metadata with a single path GET and does not list folder children" do + metadata_response = instance_double( + Net::HTTPSuccess, + code: "200", + body: { + "id" => "csv-nested", + "name" => "sales.csv", + "parentReference" => { "driveId" => "drive-1" } + }.to_json + ) + path_url = "https://graph.microsoft.com/v1.0/drives/drive-1/root:/reports/2024/sales.csv" + + expect(Multiwoven::Integrations::Core::HttpClient).to receive(:request).with( + path_url, + Multiwoven::Integrations::Core::Constants::HTTP_GET, + headers: graph_auth_headers + ).once.and_return(metadata_response) + expect(client).not_to receive(:fetch_list_items) + expect(client).not_to receive(:collect_files_from_folder) + + item = client.send(:fetch_drive_item_by_path, "reports/2024/sales.csv") + + expect(item["id"]).to eq("csv-nested") + expect(item.dig("parentReference", "driveId")).to eq("drive-1") + end + + it "returns nil when the path metadata request fails" do + allow(client).to receive(:microsoft_graph_request).and_raise(StandardError, "not found") + + expect(client.send(:fetch_drive_item_by_path, "missing.csv")).to be_nil + end + end + + describe "#resolve_download_file_name" do + it "preserves relative paths for non-http inputs" do + expect(client.send(:resolve_download_file_name, "docs/nested.pdf")).to eq("docs/nested.pdf") + end + + it "uses configured file_name when the input is an http URL" do + client.instance_variable_set(:@file_name, "configured.pdf") + + expect(client.send(:resolve_download_file_name, "https://example.com/share/file")) + .to eq("configured.pdf") + end + end + + describe "#matching_file_name?" do + it "matches by relative path when configured" do + client.instance_variable_set(:@file_name, "docs/nested.pdf") + + expect(client.send(:matching_file_name?, "nested.pdf", "docs/nested.pdf")).to be(true) + expect(client.send(:matching_file_name?, "other.pdf", "docs/other.pdf")).to be(false) + end + end + end + + describe "connector interface" do + it "exposes meta_data with the connector module name" do + meta_name = client.class.to_s.split("::")[-2] + expect(client.send(:meta_data)[:data][:name]).to eq(meta_name) + end + + it "defines private #create_connection and #query methods" do + private_methods = described_class.private_instance_methods + expect(private_methods).to include(:create_connection, :query) + end + end +end