% Copyright (C) 2026 Kocoy Group and AsaDB contributors % SPDX-License-Identifier: GPL-3.0-only /* AsaDB Local Web API + AsAPanel static host. */ :- set_prolog_flag(double_quotes, codes). :- use_module(library(http/thread_httpd)). :- use_module(library(http/http_dispatch)). :- use_module(library(http/http_parameters)). :- use_module(library(http/http_client)). :- use_module(library(http/http_multipart_plugin)). :- use_module(library(http/json)). :- use_module(library(http/http_stream)). :- use_module(library(random)). :- use_module(library(readutil)). :- use_module(library(utf8)). :- use_module('asadb_core.pl'). :- use_module('asadb_backup.pl'). :- use_module('asadb_config.pl'). :- use_module('asadb_interchange.pl'). :- use_module('bridge/reservoir.pl'). :- initialization(main, main). :- dynamic asadb_panel_token/1. :- dynamic asadb_import_progress/9. :- dynamic asadb_import_pending/2. :- http_handler(root(.), serve_index, []). :- http_handler(root('assets/app-loader.js'), serve_asset('web/assets/app-loader.js', 'application/javascript'), []). :- http_handler(root('assets/app.js'), serve_asset('web/assets/app.js', 'application/javascript'), []). :- http_handler(root('assets/app.legacy.js'), serve_asset('web/assets/app.legacy.js', 'application/javascript'), []). :- http_handler(root('assets/style.css'), serve_asset('web/assets/style.css', 'text/css'), []). :- http_handler(root('favicon.ico'), serve_binary_asset('web/assets/asadb-logo.png', 'image/png'), []). :- http_handler(root('assets/fonts/noto-sans-jp-japanese-400-normal.woff2'), serve_binary_asset('web/assets/fonts/noto-sans-jp-japanese-400-normal.woff2', 'font/woff2'), []). :- http_handler(root('assets/fonts/noto-sans-jp-japanese-400-normal.woff'), serve_binary_asset('web/assets/fonts/noto-sans-jp-japanese-400-normal.woff', 'font/woff'), []). :- http_handler(root('assets/asadb-logo.png'), serve_binary_asset('web/assets/asadb-logo.png', 'image/png'), []). :- http_handler(root('assets/Icon/save.png'), serve_binary_asset('web/assets/Icon/save.png', 'image/png'), []). :- http_handler(root('assets/Icon/tong_sampah.png'), serve_binary_asset('web/assets/Icon/tong_sampah.png', 'image/png'), []). :- http_handler(root('assets/Effect/Berhasil/1.mp3'), serve_binary_asset('web/assets/Effect/Berhasil/1.mp3', 'audio/mpeg'), []). :- http_handler(root('assets/Effect/Berhasil/2.mp3'), serve_binary_asset('web/assets/Effect/Berhasil/2.mp3', 'audio/mpeg'), []). :- http_handler(root('assets/Effect/Berhasil/3.mp3'), serve_binary_asset('web/assets/Effect/Berhasil/3.mp3', 'audio/mpeg'), []). :- http_handler(root('assets/Effect/Berhasil/4.mp3'), serve_binary_asset('web/assets/Effect/Berhasil/4.mp3', 'audio/mpeg'), []). :- http_handler(root('assets/Effect/Gagal/1g.mp3'), serve_binary_asset('web/assets/Effect/Gagal/1g.mp3', 'audio/mpeg'), []). :- http_handler(root('assets/Effect/Gagal/2g.mp3'), serve_binary_asset('web/assets/Effect/Gagal/2g.mp3', 'audio/mpeg'), []). :- http_handler(root('assets/Effect/Gagal/3g.mp3'), serve_binary_asset('web/assets/Effect/Gagal/3g.mp3', 'audio/mpeg'), []). :- http_handler(root('assets/Effect/Gagal/4g.mp3'), serve_binary_asset('web/assets/Effect/Gagal/4g.mp3', 'audio/mpeg'), []). :- http_handler(root('samples/demo.sql'), serve_asset('web/samples/demo.sql', 'application/sql'), []). :- http_handler(root('samples/feature-tour.sql'), serve_asset('web/samples/feature-tour.sql', 'application/sql'), []). :- http_handler(root('api/query'), api_query, []). :- http_handler(root('api/execute_stream'), api_execute_stream, []). :- http_handler(root('api/warmup'), api_warmup, []). :- http_handler(root('api/save'), api_save, []). :- http_handler(root('api/shutdown'), api_shutdown, []). :- http_handler(root('api/analyze'), api_analyze, []). :- http_handler(root('api/state'), api_state, []). :- http_handler(root('api/catalog'), api_catalog, []). :- http_handler(root('api/metadata'), api_metadata, []). :- http_handler(root('api/backup'), api_backup, []). :- http_handler(root('api/export'), api_export, []). :- http_handler(root('api/import_file'), api_import_file, []). :- http_handler(root('api/import_upload'), api_import_upload, []). :- http_handler(root('api/import_progress'), api_import_progress, []). :- http_handler(root('api/reservoir/jobs'), api_reservoir_jobs, []). :- http_handler(root('api/reservoir/file'), api_reservoir_file, []). :- http_handler(root('api/reservoir/job'), api_reservoir_job, []). :- http_handler(root('api/reservoir/result'), api_reservoir_result, []). :- http_handler(root('api/reservoir/cancel'), api_reservoir_cancel, []). :- http_handler(root('api/reservoir/stats'), api_reservoir_stats, []). main :- current_prolog_flag(argv, Argv), parse_args(Argv, DbFile, Port), init_panel_token, asadb_boot(DbFile), asadb_warmup, asadb_init_reservoir(DbFile), asadb_start_http_on_available_port(Port, ActualPort), asadb_write_panel_port_file(ActualPort), format('AsAPanel running at http://127.0.0.1:~w/ using ~w~n', [ActualPort, DbFile]), wait_forever. parse_args([DbFile, PortAtom|_], DbFile, Port) :- atom_number(PortAtom, Port), !. parse_args([DbFile|_], DbFile, 8088) :- !. parse_args(_, 'data.asa', 8088). wait_forever :- sleep(3600), wait_forever. asadb_start_http_on_available_port(PreferredPort, Port) :- Upper is PreferredPort + 30, between(PreferredPort, Upper, Port), catch(http_server(http_dispatch, [port(localhost:Port), silent(true)]), _, fail), !. asadb_start_http_on_available_port(PreferredPort, _) :- format(user_error, 'AsA fatal: no local port available near ~w.~n', [PreferredPort]), halt(1). asadb_write_panel_port_file(Port) :- catch( setup_call_cleanup( open('asadb.port', write, Out), format(Out, 'http://127.0.0.1:~w/~n', [Port]), close(Out) ), _, true ). init_panel_token :- retractall(asadb_panel_token(_)), random_between(100000000, 999999999, A), random_between(100000000, 999999999, B), random_between(100000000, 999999999, C), format(atom(Token), '~w-~w-~w', [A, B, C]), assertz(asadb_panel_token(Token)). asadb_init_reservoir(DbFile) :- reservoir_init(DbFile, user:reservoir_execute_job). reservoir_execute_job(JobId, SpoolPath, StopOnError, Result) :- with_mutex(asadb_execution, reservoir_execute_job_scoped(JobId, SpoolPath, StopOnError, Result)). reservoir_execute_job_scoped(JobId, SpoolPath, StopOnError, Result) :- asadb_backup_capture_current_database(PreviousDatabase), setup_call_cleanup( true, reservoir_execute_spooled_sql(JobId, SpoolPath, StopOnError, Result), restore_import_database_safely(PreviousDatabase) ). reservoir_execute_spooled_sql(JobId, SpoolPath, StopOnError, Result) :- reservoir_job_snapshot(JobId, Snapshot), get_dict(metadata, Snapshot, Metadata), reservoir_select_logical_database(Metadata), reservoir_execute_spooled_sql_metadata(Metadata, JobId, SpoolPath, StopOnError, Result). reservoir_execute_spooled_sql_metadata(Metadata, JobId, SpoolPath, StopOnError, Result) :- get_dict(kind, Metadata, interchange), !, reservoir_interchange_options(Metadata, Format, OriginalName, Target, Mode), asadb_interchange_prepare_import(SpoolPath, OriginalName, Format, Target, Mode, Prepared, _), setup_call_cleanup( true, reservoir_execute_prepared_sql(JobId, Prepared, StopOnError, Result), asadb_interchange_cleanup(Prepared) ). reservoir_execute_spooled_sql_metadata(_, JobId, SpoolPath, StopOnError, Result) :- import_sql_file_backend(SpoolPath, StopOnError, JobId, Result). reservoir_execute_prepared_sql(JobId, Prepared, StopOnError, Result) :- import_sql_file_backend(Prepared, StopOnError, JobId, Result). reservoir_interchange_options(Metadata, Format, OriginalName, Target, Mode) :- interchange_metadata_value(Metadata, format, auto, Format), interchange_metadata_value(Metadata, source_name, 'import.sql', OriginalName), interchange_metadata_value(Metadata, target_table, '', Target), interchange_metadata_value(Metadata, mode, replace, Mode). interchange_metadata_value(Dict, Key, Default, Value) :- ( get_dict(Key, Dict, Found) -> Value = Found ; Value = Default ). % Server-mode jobs are asynchronous. Capture the browser's logical database % at admission and restore it while holding the same execution mutex as BEGIN % and the import itself. A later request therefore cannot redirect a queued % CSV/XLSX/SQL job into another database. reservoir_select_logical_database(Metadata) :- get_dict(logical_database, Metadata, Database), memberchk(Database, [none, '__asadb_no_database__']), !, asadb_backup_restore_current_database(none). reservoir_select_logical_database(Metadata) :- get_dict(logical_database, Metadata, Database), atom(Database), Database \== '', !, asadb_backup_restore_current_database(Database). reservoir_select_logical_database(_). result_statement_count(multi(Results), Count) :- !, length(Results, Count). result_statement_count(_, 1). serve_index(_Request) :- asadb_panel_token(Token), security_headers, format('Set-Cookie: asadb_token=~w; Path=/; SameSite=Strict~n', [Token]), serve_file_body('web/index.html', 'text/html'). serve_asset(Path, Type, _Request) :- serve_file(Path, Type). serve_binary_asset(Path, Type, _Request) :- serve_binary_file(Path, Type). serve_file(Path, Type) :- security_headers, serve_file_body(Path, Type). serve_file_body(Path, Type) :- % Static UI files include Japanese translations. Read them as UTF-8 text % so SWI-Prolog transcodes characters once for the HTTP text stream. % Copying binary bytes to that stream double-encodes non-ASCII text. format('Content-type: ~w; charset=utf-8~n~n', [Type]), setup_call_cleanup( open(Path, read, In, [encoding(utf8)]), copy_stream_data(In, current_output), close(In) ). serve_binary_file(Path, Type) :- security_headers, size_file(Path, Size), format('Content-type: ~w~n', [Type]), format('Content-length: ~w~n~n', [Size]), setup_call_cleanup( open(Path, read, In, [type(binary)]), copy_stream_data(In, current_output), close(In) ). api_query(Request) :- member(method(post), Request), !, ( authorized_api(Request) -> ( read_sql_page_body(Request, SQL, Offset) -> web_query_fetch_size(FetchRows), web_query_execution(SQL, Offset, FetchRows, Result0, OffsetApplied), ( OffsetApplied == true -> web_query_result(Result0, Result) ; web_query_result_offset(Result0, Offset, Result) ), asadb_result_json(Result, JSON), json_response(JSON) ; json_error('400 Bad Request', 'Missing or oversized SQL payload') ) ; json_error('403 Forbidden', 'Forbidden') ). api_query(_) :- json_error('405 Method Not Allowed', 'POST only'). read_sql_page_body(Request, SQL, Offset) :- http_read_data(Request, Data, []), member(sql=SQL, Data), atom_length(SQL, Len), Len =< 250000, query_page_offset(Data, Offset). query_page_offset(Data, Offset) :- member(offset=Raw, Data), !, page_number_value(Raw, Offset). query_page_offset(_, 0). page_number_value(Value, Number) :- integer(Value), Value >= 0, Value =< 1000000000, Number = Value, !. page_number_value(Value, Number) :- atom(Value), catch(atom_number(Value, Number), _, fail), integer(Number), Number >= 0, Number =< 1000000000. query_page_snapshot_execution(SQL, Offset, FetchRows, Result, true) :- Offset > 0, asadb_exec_sql_snapshot_page(SQL, Offset, FetchRows, Result). query_page_snapshot_execution(SQL, _Offset, FetchRows, Result, false) :- asadb_exec_sql_snapshot_limited(SQL, FetchRows, Result). query_page_execution(SQL, Offset, FetchRows, Result, true) :- Offset > 0, catch(asadb_parse_sql(SQL, [select(_, _, _, _, _, _)]), _, fail), !, asadb_exec_sql_page(SQL, Offset, FetchRows, Result). query_page_execution(SQL, _Offset, FetchRows, Result, false) :- asadb_exec_sql_limited(SQL, FetchRows, Result). % SELECT receives an immutable TVCC generation and no longer waits behind a % long import or writer. Every other statement keeps the established local % single-writer execution mutex, including transactions and administrative % commands. web_query_execution(SQL, Offset, FetchRows, Result, OffsetApplied) :- asadb_snapshot_read_allowed(SQL), !, query_page_snapshot_execution(SQL, Offset, FetchRows, Result, OffsetApplied). web_query_execution(SQL, Offset, FetchRows, Result, OffsetApplied) :- with_mutex(asadb_execution, query_page_execution(SQL, Offset, FetchRows, Result, OffsetApplied)). api_execute_stream(Request) :- member(method(post), Request), !, ( authorized_api(Request) -> ( request_stream_payload(Request, In, Size) -> request_stop_on_error(Request, StopOnError), request_idempotency_key(Request, IdempotencyKey), setup_call_cleanup( true, catch( ( reservoir_submit_stream(In, 'SQL command stream', Size, IdempotencyKey, StopOnError, JobId, _), reservoir_wait(JobId, 3600, WaitOutcome), reservoir_wait_result(WaitOutcome, Result) ), Error, Result = error(stream_execution_failed, Error) ), close(In) ), asadb_result_json(Result, JSON), json_response(JSON) ; json_error('400 Bad Request', 'Missing or oversized SQL stream') ) ; json_error('403 Forbidden', 'Forbidden') ). api_execute_stream(_) :- json_error('405 Method Not Allowed', 'POST only'). reservoir_wait_result(result(Result), Result) :- !. reservoir_wait_result(error(Status, Message), error(Status, Message)) :- !. reservoir_wait_result(timeout, error(reservoir_timeout, 'Reservoir job timed out while waiting for compatibility response.')). web_query_row_limit(Limit) :- asadb_config_get(max_result_rows, Limit). web_query_fetch_size(FetchRows) :- web_query_row_limit(Limit), FetchRows is Limit + 1. web_query_result(multi(Results0), multi(Results)) :- !, maplist(web_query_result, Results0, Results). web_query_result(table(Columns, Rows0), table_page(Columns, Rows, HasMore)) :- !, web_query_row_limit(Limit), take_web_rows(Limit, Rows0, Rows, HasMore). web_query_result(Result, Result). web_query_result_offset(multi(Results0), Offset, multi(Results)) :- !, maplist(web_query_result_offset_at(Offset), Results0, Results). web_query_result_offset(table(Columns, Rows0), Offset, table_page(Columns, Rows, HasMore)) :- !, web_query_row_limit(Limit), drop_web_rows(Offset, Rows0, Remaining), take_web_rows(Limit, Remaining, Rows, HasMore). web_query_result_offset(Result, _, Result). web_query_result_offset_at(Offset, Result0, Result) :- web_query_result_offset(Result0, Offset, Result). drop_web_rows(0, Rows, Rows) :- !. drop_web_rows(_, [], []) :- !. drop_web_rows(N, [_|Rows], Remaining) :- N1 is max(0, N - 1), drop_web_rows(N1, Rows, Remaining). take_web_rows(0, Rows, [], HasMore) :- !, ( Rows == [] -> HasMore = false ; HasMore = true ). take_web_rows(_, [], [], false) :- !. take_web_rows(N, [Row|Rows0], [Row|Rows], HasMore) :- N1 is N - 1, take_web_rows(N1, Rows0, Rows, HasMore). api_warmup(Request) :- ( authorized_api(Request) -> asadb_warmup, asadb_result_json(ok(warmup_ready), JSON), json_response(JSON) ; json_error('403 Forbidden', 'Forbidden') ). api_save(Request) :- member(method(post), Request), !, ( authorized_api(Request) -> with_mutex(asadb_execution, ( asadb_save, asadb_current_database(CurrentDb) )), asadb_result_json(ok(saved_database(CurrentDb)), JSON), json_response(JSON) ; json_error('403 Forbidden', 'Forbidden') ). api_save(_) :- json_error('405 Method Not Allowed', 'POST only'). % The supervised Flask server uses this authenticated localhost endpoint for % a normal restart. On Windows, terminating the SWI process is equivalent to % a forced kill and does not run Prolog cleanup hooks. Reply first, then stop % the runtime from a detached thread so the HTTP response reaches the % supervisor before halt/0 closes the listener. api_shutdown(Request) :- member(method(post), Request), !, ( authorized_api(Request) -> with_mutex(asadb_execution, asadb_shutdown), thread_create(( sleep(0.25), halt(0) ), _, [detached(true)]), asadb_result_json(ok(shutting_down), JSON), json_response(JSON) ; json_error('403 Forbidden', 'Forbidden') ). api_shutdown(_) :- json_error('405 Method Not Allowed', 'POST only'). api_analyze(Request) :- member(method(post), Request), !, ( authorized_api(Request) -> ( read_sql_body(Request, SQL) -> asadb_analyze_sql(SQL, Diagnostics), asadb_analysis_json(Diagnostics, JSON), json_response(JSON) ; json_error('400 Bad Request', 'Missing or oversized SQL payload') ) ; json_error('403 Forbidden', 'Forbidden') ). api_analyze(_) :- json_error('405 Method Not Allowed', 'POST only'). api_state(Request) :- ( authorized_api(Request) -> asadb_get_state(State), asadb_current_database(CurrentDb), asadb_result_json(table([state,current_db], [[State,CurrentDb]]), JSON), json_response(JSON) ; json_error('403 Forbidden', 'Forbidden') ). api_catalog(Request) :- ( authorized_api(Request) -> asadb_current_database(CurrentDb), asadb_get_state(state(_, DBs)), findall(Row, catalog_row(CurrentDb, DBs, Row), Rows), asadb_result_json(table([current_db,database,kind,name,row_count,columns,indexes,query], Rows), JSON), json_response(JSON) ; json_error('403 Forbidden', 'Forbidden') ). api_metadata(Request) :- ( authorized_api(Request) -> asadb_database_metadata(CoreMetadata), reservoir_stats(Reservoir), Metadata = CoreMetadata.put(reservoir, Reservoir), json_dict_response(Metadata) ; json_error('403 Forbidden', 'Forbidden') ). % The backup endpoint intentionally accepts only database/output controls. % It never accepts table rows, browser state, or a client-side catalog: the % production artifact is generated by scanning the backend record store. api_backup(Request) :- member(method(post), Request), !, ( authorized_api(Request) -> catch(api_backup_authorized(Request), Error, api_backup_error(Error)) ; json_error('403 Forbidden', 'Forbidden') ). api_backup(_) :- json_error('405 Method Not Allowed', 'POST only'). api_backup_authorized(Request) :- http_read_data(Request, Data, []), backup_request_database(Data, Database), backup_request_output(Data, Output), asadb_backup_create(Database, BackupFile, Manifest), setup_call_cleanup( true, serve_production_backup(BackupFile, Database, Output, Manifest), asadb_backup_cleanup(BackupFile) ). backup_request_database(Data, Database) :- member(database=Raw, Data), backup_database_value(Raw, Database), !. backup_request_database(_, _) :- throw(error(domain_error(backup_database, missing), _)). backup_database_value(Value, Database) :- atom(Value), !, Database = Value. backup_database_value(Value, Database) :- string(Value), !, atom_string(Database, Value). backup_database_value(Value, _) :- throw(error(domain_error(backup_database, Value), _)). backup_request_output(Data, Output) :- member(output=Raw, Data), member(Raw, [save,open]), !, Output = Raw. backup_request_output(_, save). serve_production_backup(File, Database, Output, Manifest) :- size_file(File, Size), backup_download_name(Database, Name), security_headers, format('Content-type: application/octet-stream~n'), format('Content-length: ~w~n', [Size]), format('X-AsaDB-Backup-Format: ~w~n', [Manifest.format]), format('X-AsaDB-Backup-SHA256: ~w~n', [Manifest.payload_sha256]), ( Output == open -> Disposition = inline ; Disposition = attachment ), format('Content-Disposition: ~w; filename="~w"~n~n', [Disposition, Name]), setup_call_cleanup( open(File, read, In, [type(binary)]), copy_stream_data(In, current_output), close(In) ). backup_download_name(Database, Name) :- atom_codes(Database, Codes), maplist(backup_filename_code, Codes, SafeCodes), atom_codes(Safe, SafeCodes), atomic_list_concat([Safe, '.asb'], Name). backup_filename_code(Code, Code) :- ( Code >= 65, Code =< 90 ; Code >= 97, Code =< 122 ; Code >= 48, Code =< 57 ; memberchk(Code, [45,95]) ), !. backup_filename_code(_, 95). api_backup_error(Error) :- term_atom_safe(Error, Message), json_error('400 Bad Request', Message). % Portable exports are separate from authenticated .asb backups. Like the % backup endpoint, this endpoint accepts selection controls only; all rows are % scanned from backend storage by the Prolog interchange module. api_export(Request) :- member(method(post), Request), !, ( authorized_api(Request) -> catch(api_export_authorized(Request), Error, api_backup_error(Error)) ; json_error('403 Forbidden', 'Forbidden') ). api_export(_) :- json_error('405 Method Not Allowed', 'POST only'). api_export_authorized(Request) :- http_read_data(Request, Data, []), backup_request_database(Data, Database), backup_request_output(Data, Output), export_request_format(Data, Format), export_request_options(Data, Options), asadb_interchange_export(Database, Format, Options, File, Metadata), setup_call_cleanup( true, serve_interchange_export(File, Output, Metadata), asadb_interchange_cleanup(File) ). export_request_format(Data, Format) :- member(format=Raw, Data), memberchk(Raw, [mysql,postgresql,csv,xlsx]), !, Format = Raw. export_request_format(_, _) :- throw(error(domain_error(interchange_export_format, missing), _)). export_request_options(Data, Options) :- export_form_tables(Data, Tables), export_form_named_tables(Data, data_tables, all, DataTables), export_form_bool(Data, include_schema, true, IncludeSchema), export_form_bool(Data, include_data, true, IncludeData), export_form_bool(Data, create_database, true, CreateDatabase), export_form_bool(Data, drop_tables, true, DropTables), Options = _{ tables:Tables, data_tables:DataTables, include_schema:IncludeSchema, include_data:IncludeData, create_database:CreateDatabase, drop_tables:DropTables }. export_form_tables(Data, Tables) :- member(tables=Raw, Data), atom(Raw), Raw \== '', !, atomic_list_concat(Parts0, ',', Raw), exclude(=(''), Parts0, Tables). export_form_tables(_, all). export_form_named_tables(Data, Key, Default, Tables) :- ( member(Pair, Data), Pair = (FoundKey=Raw), FoundKey == Key, atom(Raw) -> ( Raw == '' -> Tables = [] ; atomic_list_concat(Parts0, ',', Raw), exclude(=(''), Parts0, Tables) ) ; Tables = Default ). export_form_bool(Data, Name, Default, Value) :- ( member(Pair, Data), Pair = (FoundName=Raw), FoundName == Name -> ( memberchk(Raw, [true,'true',yes,'yes','1',1,on,'on']) -> Value = true ; Value = false ) ; Value = Default ). serve_interchange_export(File, Output, Metadata) :- security_headers, format('Content-type: ~w~n', [Metadata.content_type]), format('Content-length: ~w~n', [Metadata.bytes]), format('X-AsaDB-Export-Format: ~w~n', [Metadata.format]), format('X-AsaDB-Export-Rows: ~w~n', [Metadata.row_count]), ( Output == open -> Disposition = inline ; Disposition = attachment ), format('Content-Disposition: ~w; filename="~w"~n~n', [Disposition, Metadata.filename]), setup_call_cleanup( open(File, read, In, [type(binary)]), copy_stream_data(In, current_output), close(In) ). api_reservoir_jobs(Request) :- member(method(get), Request), !, ( authorized_api(Request) -> reservoir_job_snapshots(Snapshots), length(Snapshots, Count), json_dict_response(reservoir_jobs{count:Count,jobs:Snapshots}) ; json_error('403 Forbidden', 'Forbidden') ). api_reservoir_jobs(Request) :- member(method(post), Request), !, ( authorized_api(Request) -> ( request_stream_payload(Request, In, Size) -> request_stop_on_error(Request, StopOnError), request_idempotency_key(Request, IdempotencyKey), request_job_label(Request, Label), request_reservoir_metadata(Request, Metadata), setup_call_cleanup( true, catch( ( reservoir_submit_stream(In, Label, Size, IdempotencyKey, StopOnError, JobId, Admission, Metadata), reservoir_job_snapshot(JobId, Snapshot), Response = reservoir_admission{ status:accepted, admission:Admission, job_id:JobId, job:Snapshot }, json_dict_response('202 Accepted', Response) ), Error, reservoir_api_error(Error) ), close(In) ) ; json_error('400 Bad Request', 'Missing or oversized Reservoir payload') ) ; json_error('403 Forbidden', 'Forbidden') ). api_reservoir_jobs(_) :- json_error('405 Method Not Allowed', 'GET or POST only'). api_reservoir_file(Request) :- member(method(post), Request), !, ( authorized_api(Request) -> http_read_data(Request, Data, []), ( member(path=RawPath, Data), allowed_import_path(RawPath, File), exists_file(File) -> catch( ( import_stop_on_error(Data, StopOnError), form_idempotency_key(Data, IdempotencyKey), import_interchange_options(Data, Format, Target, Mode), request_reservoir_database(Request, Database), Metadata = _{ kind:interchange, format:Format, source_name:File, target_table:Target, mode:Mode, logical_database:Database }, reservoir_submit_file(File, File, IdempotencyKey, StopOnError, JobId, Admission, Size, Metadata), reservoir_job_snapshot(JobId, Snapshot), Response = reservoir_admission{ status:accepted, admission:Admission, job_id:JobId, size_bytes:Size, job:Snapshot }, json_dict_response('202 Accepted', Response) ), Error, reservoir_api_error(Error) ) ; json_error('400 Bad Request', 'Import file path is not allowed or not found') ) ; json_error('403 Forbidden', 'Forbidden') ). api_reservoir_file(_) :- json_error('405 Method Not Allowed', 'POST only'). api_reservoir_job(Request) :- member(method(get), Request), !, ( authorized_api(Request) -> catch( ( reservoir_http_job_id(Request, JobId), reservoir_job_snapshot(JobId, Snapshot), json_dict_response(Snapshot) ), Error, reservoir_api_error(Error) ) ; json_error('403 Forbidden', 'Forbidden') ). api_reservoir_job(_) :- json_error('405 Method Not Allowed', 'GET only'). api_reservoir_result(Request) :- member(method(get), Request), !, ( authorized_api(Request) -> catch( ( reservoir_http_result_parameters(Request, JobId, Offset, Limit), reservoir_job_result_page(JobId, Offset, Limit, Result), asadb_result_json(Result, JSON), json_response(JSON) ), Error, reservoir_api_error(Error) ) ; json_error('403 Forbidden', 'Forbidden') ). api_reservoir_result(_) :- json_error('405 Method Not Allowed', 'GET only'). api_reservoir_cancel(Request) :- member(method(post), Request), !, ( authorized_api(Request) -> catch( ( reservoir_http_job_id(Request, JobId), reservoir_cancel(JobId, Snapshot), json_dict_response(Snapshot) ), Error, reservoir_api_error(Error) ) ; json_error('403 Forbidden', 'Forbidden') ). api_reservoir_cancel(_) :- json_error('405 Method Not Allowed', 'POST only'). api_reservoir_stats(Request) :- member(method(get), Request), !, ( authorized_api(Request) -> reservoir_stats(Stats), json_dict_response(Stats) ; json_error('403 Forbidden', 'Forbidden') ). api_reservoir_stats(_) :- json_error('405 Method Not Allowed', 'GET only'). api_import_file(Request) :- member(method(post), Request), !, ( authorized_api(Request) -> catch(api_import_file_authorized(Request), Error, api_import_upload_error(Error)) ; json_error('403 Forbidden', 'Forbidden') ). api_import_file(_) :- json_error('405 Method Not Allowed', 'POST only'). api_import_file_authorized(Request) :- http_read_data(Request, Data, []), ( member(path=RawPath, Data), allowed_import_path(RawPath, File), exists_file(File) -> import_stop_on_error(Data, StopOnError), import_id(Data, ImportId), import_interchange_options(Data, Format, Target, Mode), with_mutex(asadb_execution, import_uploaded_database_file(File, File, Format, Target, Mode, StopOnError, ImportId, Result)), asadb_result_json(Result, JSON), json_response(JSON) ; json_error('400 Bad Request', 'Import file path is not allowed or not found') ). api_import_upload(Request) :- member(method(post), Request), !, ( authorized_api(Request) -> catch(api_import_upload_authorized(Request), Error, api_import_upload_error(Error)) ; json_error('403 Forbidden', 'Forbidden') ). api_import_upload(_) :- json_error('405 Method Not Allowed', 'POST only'). api_import_upload_authorized(Request) :- http_read_data(Request, Data, [on_filename(save_uploaded_import_part)]), ( member(file=upload(TempFile, OriginalName, Size), Data) -> import_stop_on_error(Data, StopOnError), import_id(Data, ImportId), import_interchange_options(Data, Format, Target, Mode), setup_call_cleanup( true, with_mutex(asadb_execution, import_uploaded_database_file(TempFile, OriginalName, Format, Target, Mode, StopOnError, ImportId, Result0)), cleanup_uploaded_import_file(TempFile) ), uploaded_import_result(Result0, OriginalName, Size, Result), asadb_result_json(Result, JSON), json_response(JSON) ; json_error('400 Bad Request', 'No SQL upload file found') ). % A production backup is verified before it reaches the SQL importer. Its % catalog-only objects are staged immediately before COMMIT, making the full % restore atomic: SQL rows and catalog objects succeed or roll back together. import_uploaded_database_file(File, OriginalName, Format, Target, Mode, StopOnError, ImportId, Result) :- ( production_backup_candidate(File, OriginalName) -> asadb_backup_prepare_restore(File, Database, PayloadFile, Manifest), setup_call_cleanup( true, import_verified_production_backup(PayloadFile, Database, Manifest, StopOnError, ImportId, Result), asadb_backup_cleanup(PayloadFile) ) ; asadb_interchange_prepare_import(File, OriginalName, Format, Target, Mode, Prepared, _), setup_call_cleanup( true, import_sql_file_backend(Prepared, StopOnError, ImportId, Result), asadb_interchange_cleanup(Prepared) ) ). import_interchange_options(Data, Format, Target, Mode) :- import_form_value(Data, format, auto, Format), import_form_value(Data, target_table, '', Target), import_form_value(Data, mode, replace, Mode). import_form_value(Data, Key, Default, Value) :- ( member(Pair, Data), Pair = (FoundKey=Found), FoundKey == Key, Found \== '' -> Value = Found ; Value = Default ). % A file called .asb is never permitted to fall through into the generic SQL % importer. That prevents a damaged production backup from being interpreted % as comment-prefixed SQL and bypassing its integrity checks. The marker % probe also catches the common "one broken byte in the magic" case when a % file has been renamed before upload. production_backup_candidate(_File, OriginalName) :- production_backup_name(OriginalName), !. production_backup_candidate(File, _) :- asadb_backup_file(File), !. production_backup_candidate(File, _) :- production_backup_marker_present(File). production_backup_name(Name) :- atom(Name), file_name_extension(_, Extension0, Name), downcase_atom(Extension0, asb), !. production_backup_name(Name) :- string(Name), atom_string(Atom, Name), production_backup_name(Atom). production_backup_marker_present(File) :- setup_call_cleanup( open(File, read, In, [type(binary)]), ( read_string(In, 4096, Prefix), sub_string(Prefix, _, _, _, 'ASADB-PRODUCTION-BACKUP') ), close(In) ). import_verified_production_backup(PayloadFile, Database, Manifest, StopOnError, ImportId, Result) :- CatalogObjects = Manifest.catalog_objects, CatalogObjects = catalog_objects(Views, Functions, Procedures, Triggers), import_sql_file_backend_before_commit( PayloadFile, StopOnError, ImportId, import_verified_backup_precommit(Database, Manifest, Views, Functions, Procedures, Triggers), Result ). import_verified_backup_precommit(Database, Manifest, Views, Functions, Procedures, Triggers) :- asadb_backup_validate_restored_snapshot(Database, Manifest), asadb_backup_stage_catalog_objects(Database, Views, Functions, Procedures, Triggers). api_import_upload_error(Error) :- term_atom_safe(Error, Message), json_error('400 Bad Request', Message). term_atom_safe(Term, Atom) :- catch(message_to_string(Term, String), _, fail), !, atom_string(Atom, String). term_atom_safe(Term, Atom) :- with_output_to(atom(Atom), write_term(Term, [quoted(true), max_depth(8)])). save_uploaded_import_part(In, upload(TempFile, BaseName, Size), Options) :- memberchk(filename(ClientName), Options), uploaded_file_base(ClientName, BaseName), safe_import_file_name(BaseName), !, tmp_file_stream(octet, TempFile, Out), setup_call_cleanup( true, copy_stream_data(In, Out), close(Out) ), size_file(TempFile, Size). save_uploaded_import_part(_, _, Options) :- memberchk(filename(ClientName), Options), uploaded_file_base(ClientName, BaseName), throw(error(permission_error(import, file, BaseName), _)). uploaded_file_base(Raw, Base) :- atom_codes(Raw, Codes0), maplist(import_slash_code, Codes0, Codes), atom_codes(Slashed, Codes), atomic_list_concat(Parts0, '/', Slashed), exclude(=(''), Parts0, Parts), Parts \= [], last(Parts, Base). cleanup_uploaded_import_file(File) :- ( exists_file(File) -> catch(delete_file(File), _, true) ; true ). uploaded_import_result(table(Columns, [[Status, _TempPath, Statements, Errors, LastStatus, LastMessage, _TempSize]]), OriginalName, Size, table(Columns, [[Status,OriginalName,Statements,Errors,LastStatus,LastMessage,Size]])) :- !. uploaded_import_result(Result, _, _, Result). api_import_progress(Request) :- ( authorized_api(Request) -> http_parameters(Request, [id(Id, [atom, default('')])]), import_progress_result(Id, Result), asadb_result_json(Result, JSON), json_response(JSON) ; json_error('403 Forbidden', 'Forbidden') ). allowed_import_path(RawPath, File) :- normalize_import_path(RawPath, Clean), allowed_import_path_clean(Clean, File). allowed_import_path_clean(Clean, File) :- atomic_list_concat(['test', Name], '/', Clean), safe_import_file_name(Name), !, stress_import_file(Name, File). allowed_import_path_clean(Clean, File) :- atomic_list_concat(['stress test', Name], '/', Clean), safe_import_file_name(Name), !, stress_import_file(Name, File). allowed_import_path_clean(Clean, File) :- atomic_list_concat(['stress-test', Name], '/', Clean), safe_import_file_name(Name), !, stress_import_file(Name, File). allowed_import_path_clean(Clean, File) :- atomic_list_concat(['stress_test', Name], '/', Clean), safe_import_file_name(Name), !, stress_import_file(Name, File). allowed_import_path_clean(Clean, File) :- atomic_list_concat(['stress tests', Name], '/', Clean), safe_import_file_name(Name), !, stress_import_file(Name, File). allowed_import_path_clean(Clean, File) :- atomic_list_concat(['samples', Name], '/', Clean), safe_import_file_name(Name), !, atom_concat('web/samples/', Name, File). allowed_import_path_clean(Clean, File) :- atomic_list_concat(['web', 'samples', Name], '/', Clean), safe_import_file_name(Name), !, atom_concat('web/samples/', Name, File). safe_import_file_name(Name) :- \+ sub_atom(Name, _, _, _, '/'), file_name_extension(_, Ext, Name), downcase_atom(Ext, Lower), member(Lower, [asb,sql,mysql,pgsql,psql,postgres,csv,zip,xlsx]). stress_import_file(Name, File) :- directory_file_path('stress tests', Name, Candidate), ( exists_file(Candidate) -> File = Candidate ; stress_import_file_casefold(Name, File) ). stress_import_file_casefold(Name, File) :- directory_files('stress tests', Entries), member(Entry, Entries), downcase_atom(Entry, Name), safe_import_file_name(Entry), directory_file_path('stress tests', Entry, File), exists_file(File), !. normalize_import_path(RawPath, Clean) :- atom_codes(RawPath, Codes0), maplist(import_slash_code, Codes0, Codes1), atom_codes(Slashed, Codes1), atomic_list_concat(Parts0, '/', Slashed), exclude(=(''), Parts0, Parts), \+ member('..', Parts), atomic_list_concat(Parts, '/', Joined), downcase_atom(Joined, Clean). import_slash_code(92, 47) :- !. import_slash_code(Code, Code). import_stop_on_error(Data, true) :- member(stop_on_error=Value, Data), member(Value, [true, 'true', yes, 'yes', '1', 1]), !. import_stop_on_error(_, false). import_id(Data, Id) :- member(import_id=Raw, Data), atom(Raw), Raw \= '', !, Id = Raw. import_id(_, Id) :- random_between(100000000, 999999999, N), format(atom(Id), 'import-~w', [N]). import_sql_file_backend(File, StopOnError, ImportId, Result) :- import_sql_file_backend_before_commit(File, StopOnError, ImportId, true, Result). import_sql_file_backend_before_commit(File, StopOnError, ImportId, BeforeCommit, Result) :- size_file(File, Size), setup_call_cleanup( open(File, read, In, [type(binary)]), import_sql_stream_backend_before_commit(In, File, Size, StopOnError, ImportId, BeforeCommit, Result), close(In) ). import_sql_stream_backend(In, Label, Size, StopOnError, ImportId, Result) :- import_sql_stream_backend_before_commit(In, Label, Size, StopOnError, ImportId, true, Result). import_sql_stream_backend_before_commit(In, Label, Size, StopOnError, ImportId, BeforeCommit, Result) :- flag(asadb_import_batches, _, 0), ensure_import_not_cancelled(ImportId), set_import_progress(ImportId, Label, Size, 0, 0, 0, running, 'starting', false), ensure_import_transaction_idle, asadb_backup_capture_current_database(PreviousDatabase), import_transaction_command('BEGIN;'), catch( import_sql_stream_backend_transaction(In, Label, Size, StopOnError, ImportId, BeforeCommit, PreviousDatabase, Result), Error, ( import_rollback_safely, restore_import_database_safely(PreviousDatabase), record_import_failure(ImportId, Label, Size, Error), throw(Error) ) ). import_sql_stream_backend_transaction(In, Label, Size, StopOnError, ImportId, BeforeCommit, PreviousDatabase, Result) :- statistics(walltime, [StreamStart|_]), setup_call_cleanup( retractall(asadb_import_pending(ImportId, _)), import_sql_stream(In, StopOnError, Label, Size, ImportId, 0, 0, 0, none, [], [], none, '', Stats), retractall(asadb_import_pending(ImportId, _)) ), statistics(walltime, [StreamEnd|_]), garbage_collect, trim_stacks, Stats = import_stats(Statements, Errors, LastStatus, LastMessage, BytesRead), ( Errors =:= 0 -> ensure_import_not_cancelled(ImportId), must_run_before_commit(BeforeCommit), statistics(walltime, [CommitStart|_]), import_transaction_command('COMMIT;'), statistics(walltime, [CommitEnd|_]), StreamMs is StreamEnd - StreamStart, CommitMs is CommitEnd - CommitStart, format(user_error, 'AsA import timing: stream=~w ms, commit=~w ms, statements=~w~n', [StreamMs, CommitMs, Statements]), set_import_progress(ImportId, Label, Size, BytesRead, Statements, Errors, committed, LastMessage, true), Result = table([status,path,statements,errors,last_status,last_message,size_bytes], [[ok,Label,Statements,Errors,LastStatus,LastMessage,Size]]) ; import_transaction_command('ROLLBACK;'), asadb_backup_restore_current_database(PreviousDatabase), set_import_progress(ImportId, Label, Size, BytesRead, Statements, Errors, rolled_back, LastMessage, true), Result = table([status,path,statements,errors,last_status,last_message,size_bytes], [[rolled_back,Label,Statements,Errors,LastStatus,LastMessage,Size]]) ). ensure_import_transaction_idle :- ( asadb_backup_transaction_active -> throw(error(permission_error(import, database, active_transaction), _)) ; true ). must_run_before_commit(Goal) :- ( call(Goal) -> true ; throw(error(integrity_error(backup_precommit_failed), _)) ). import_transaction_command(SQL) :- asadb_exec_sql(SQL, Result), ( import_transaction_result_ok(Result) -> true ; throw(error(transaction_command_failed(SQL, Result), _)) ). import_transaction_result_ok(multi(Results)) :- Results \= [], \+ member(error(_, _), Results), !. import_transaction_result_ok(ok(_)). import_rollback_safely :- catch(import_transaction_command('ROLLBACK;'), _, true). restore_import_database_safely(Database) :- catch(asadb_backup_restore_current_database(Database), _, catch(asadb_backup_restore_current_database(none), _, true)). import_sql_stream(In, StopOnError, File, Size, ImportId, Bytes0, Statements0, Errors0, Mode0, RevAcc0, Utf8Carry0, LastStatus0, LastMessage0, Stats) :- ensure_import_not_cancelled(ImportId), import_read_block(In, Size, Bytes0, BlockString), string_length(BlockString, BlockLength), ( BlockLength =:= 0 -> ( Utf8Carry0 == [] -> true ; throw(error(invalid_utf8_tail(Utf8Carry0), _)) ), import_flush_pending(RevAcc0, StopOnError, File, Size, ImportId, Bytes0, Statements0, Errors0, LastStatus0, LastMessage0, Stats) ; string_codes(BlockString, RawBytes), length(RawBytes, BlockBytes), append(Utf8Carry0, RawBytes, Utf8Bytes), phrase(utf8_codes(BlockCodes), Utf8Bytes, Utf8Carry), Bytes1 is Bytes0 + BlockBytes, scan_sql_line(BlockCodes, Mode0, RevAcc0, [], Mode, RevAcc0Out, StatementCodes0), drain_import_partial_insert(Mode, RevAcc0Out, StatementCodes0, RevAcc, StatementCodes), % The parser cap is a hard admission boundary, not a later flush hint. % A complete read block may hold several statements, so enqueue % them one-by-one and flush *before* one would cross the byte budget. enqueue_import_statements_bounded(ImportId, StatementCodes, StopOnError, Statements0, Errors0, Statements1, Errors1, LastStatus, LastMessage, Stop), % Publish the read position before a large bounded batch is executed. % Generated stress files can contain thousands of rows in one batch; % reporting only after execution made a live panel look frozen. maybe_set_import_progress(ImportId, File, Size, Bytes0, Bytes1, Statements0, Errors0, none, 'buffered; executing batch'), merge_import_last(LastStatus0, LastMessage0, LastStatus, LastMessage, LastStatus1, LastMessage1), maybe_set_import_progress(ImportId, File, Size, Bytes0, Bytes1, Statements1, Errors1, LastStatus, LastMessage1), ( Stop == true -> Stats = import_stats(Statements1, Errors1, LastStatus1, LastMessage1, Bytes1) ; import_sql_stream(In, StopOnError, File, Size, ImportId, Bytes1, Statements1, Errors1, Mode, RevAcc, Utf8Carry, LastStatus1, LastMessage1, Stats) ) ). % scan_sql_line/7 is deliberately simple and recursive. Keep its source % slice below the worker-safe envelope: a 256 KiB slice can itself become a % high live-stack allocation before the later SQL batch admission check runs. % Large multi-row INSERTs continue through the row-boundary splitter below. import_read_block_size(32768). import_read_block(_, Size, BytesRead, "") :- Size > 0, BytesRead >= Size, !. import_read_block(In, Size, BytesRead, BlockString) :- import_read_block_size(BlockSize), ( Size > 0 -> Remaining is Size - BytesRead, ReadSize is min(BlockSize, Remaining) ; ReadSize = BlockSize ), read_string(In, ReadSize, BlockString). import_flush_pending(RevAcc, StopOnError, File, Size, ImportId, BytesRead, Statements0, Errors0, LastStatus0, LastMessage0, Stats) :- reverse(RevAcc, Codes), trim_sql_codes(Codes, Trimmed), enqueue_import_statements_bounded(ImportId, [Trimmed], StopOnError, Statements0, Errors0, QueuedStatements, QueuedErrors, QueuedStatus, QueuedMessage, _), flush_import_statements(ImportId, true, StopOnError, QueuedStatements, QueuedErrors, Statements, Errors, FlushedStatus, FlushedMessage, _), merge_import_last(QueuedStatus, QueuedMessage, FlushedStatus, FlushedMessage, LastStatus, LastMessage), merge_import_last(LastStatus0, LastMessage0, LastStatus, LastMessage, FinalStatus, FinalMessage), set_import_progress(ImportId, File, Size, BytesRead, Statements, Errors, running, FinalMessage, false), Stats = import_stats(Statements, Errors, FinalStatus, FinalMessage, BytesRead). maybe_set_import_progress(ImportId, File, Size, _, Bytes, Statements, Errors, LastStatus, Message) :- LastStatus \== none, !, set_import_progress(ImportId, File, Size, Bytes, Statements, Errors, running, Message, false). maybe_set_import_progress(ImportId, File, Size, Bytes0, Bytes, Statements, Errors, _, Message) :- ProgressQuantum is 64 * 1024, PreviousBucket is Bytes0 // ProgressQuantum, CurrentBucket is Bytes // ProgressQuantum, CurrentBucket > PreviousBucket, !, set_import_progress(ImportId, File, Size, Bytes, Statements, Errors, running, Message, false). maybe_set_import_progress(_, _, _, _, _, _, _, _, _). import_statement_batch_size(BatchSize) :- asadb_config_get(import_batch_size, BatchSize). % Queue each complete statement under the exact source-size budget passed to % the SQL frontend. The previous implementation queued an entire read block % then noticed it was too large; that made 512 KiB a soft limit and could hand % approximately 733 KiB to scan/2 in a 64 MiB Reservoir worker. % The [] and [Head|Tail] clauses are disjoint, but SWI-Prolog 9.2 can retain a % choicepoint without these green cuts. That choicepoint keeps every prior % SQL code list live across batches and defeats the explicit GC boundary. enqueue_import_statements_bounded(_, [], _, Statements, Errors, Statements, Errors, none, '', false) :- !. enqueue_import_statements_bounded(ImportId, [Codes|Rest], StopOnError, Statements0, Errors0, Statements, Errors, LastStatus, LastMessage, Stop) :- !, ( Codes == [] -> enqueue_import_statements_bounded(ImportId, Rest, StopOnError, Statements0, Errors0, Statements, Errors, LastStatus, LastMessage, Stop) ; import_statement_chunks(Codes, Chunks), enqueue_import_chunks_bounded(ImportId, Chunks, StopOnError, Statements0, Errors0, Statements1, Errors1, ChunkStatus, ChunkMessage, ChunkStop), ( ChunkStop == true -> Statements = Statements1, Errors = Errors1, LastStatus = ChunkStatus, LastMessage = ChunkMessage, Stop = true ; enqueue_import_statements_bounded(ImportId, Rest, StopOnError, Statements1, Errors1, Statements, Errors, RestStatus, RestMessage, Stop), merge_import_last(ChunkStatus, ChunkMessage, RestStatus, RestMessage, LastStatus, LastMessage) ) ). enqueue_import_chunks_bounded(_, [], _, Statements, Errors, Statements, Errors, none, '', false) :- !. enqueue_import_chunks_bounded(ImportId, [Codes|Rest], StopOnError, Statements0, Errors0, Statements, Errors, LastStatus, LastMessage, Stop) :- !, import_statement_sql_bytes(Codes, StatementBytes), flush_before_import_enqueue(ImportId, StatementBytes, StopOnError, Statements0, Errors0, BeforeStatements, BeforeErrors, BeforeStatus, BeforeMessage, BeforeStop), ( BeforeStop == true -> Statements = BeforeStatements, Errors = BeforeErrors, LastStatus = BeforeStatus, LastMessage = BeforeMessage, Stop = true ; assertz(asadb_import_pending(ImportId, Codes)), flush_import_statements(ImportId, false, StopOnError, BeforeStatements, BeforeErrors, AfterStatements, AfterErrors, AfterStatus, AfterMessage, AfterStop), ( AfterStop == true -> Statements = AfterStatements, Errors = AfterErrors, LastStatus = AfterStatus, LastMessage = AfterMessage, Stop = true ; enqueue_import_chunks_bounded(ImportId, Rest, StopOnError, AfterStatements, AfterErrors, Statements, Errors, RestStatus, RestMessage, Stop), merge_import_last(BeforeStatus, BeforeMessage, AfterStatus, AfterMessage, InterimStatus, InterimMessage), merge_import_last(InterimStatus, InterimMessage, RestStatus, RestMessage, LastStatus, LastMessage) ) ). flush_before_import_enqueue(ImportId, StatementBytes, StopOnError, Statements0, Errors0, Statements, Errors, LastStatus, LastMessage, Stop) :- pending_import_metrics(ImportId, PendingCount, PendingBytes), import_statement_batch_size(BatchSize), import_execution_batch_bytes(MaxBatchBytes), ( PendingCount > 0, ( PendingCount >= BatchSize ; PendingBytes + StatementBytes > MaxBatchBytes ) -> flush_import_statements(ImportId, true, StopOnError, Statements0, Errors0, Statements, Errors, LastStatus, LastMessage, Stop) ; Statements = Statements0, Errors = Errors0, LastStatus = none, LastMessage = '', Stop = false ). pending_import_metrics(ImportId, PendingCount, PendingBytes) :- aggregate_all(count, asadb_import_pending(ImportId, _), PendingCount), aggregate_all(sum(Bytes), ( asadb_import_pending(ImportId, Codes), import_statement_sql_bytes(Codes, Bytes) ), PendingBytes). flush_import_statements(ImportId, Force, StopOnError, Statements0, Errors0, Statements, Errors, LastStatus, LastMessage, Stop) :- pending_import_metrics(ImportId, PendingCount, PendingBytes), import_statement_batch_size(BatchSize), ( PendingCount =:= 0 -> Statements = Statements0, Errors = Errors0, LastStatus = none, LastMessage = '', Stop = false ; ( Force == true ; PendingCount >= BatchSize ; import_execution_batch_bytes(MaxBatchBytes), PendingBytes >= MaxBatchBytes ) -> findall(Codes, retract(asadb_import_pending(ImportId, Codes)), CodeGroups), statement_groups_sql(CodeGroups, SQLCodes), import_execute_sql_batch(StopOnError, SQLCodes, Result), length(CodeGroups, ExecutedCount), Statements is Statements0 + ExecutedCount, result_error_count(Result, BatchErrors), Errors is Errors0 + BatchErrors, result_status_message(Result, LastStatus, LastMessage), ( StopOnError == true, BatchErrors > 0 -> Stop = true ; Stop = false ), release_import_batch_memory ; Statements = Statements0, Errors = Errors0, LastStatus = none, LastMessage = '', Stop = false ). % Statement count alone is not a memory bound: generated dumps often contain % 50-100 KiB multi-row INSERT statements. Cap the source text handed to one % parser run while preserving the higher count limit for tiny statements. import_statement_batch_bytes(524288). % 512 KiB remains the documented parser ceiling. A Reservoir worker also % keeps a deliberately smaller execution envelope: parsing several narrow % 512-row INSERTs together can expand into thousands of row terms before the % storage writer runs. This is a stricter cap, never a relaxation of the % configured parser maximum. import_execution_batch_bytes(Bytes) :- import_statement_batch_bytes(ParserBytes), import_worker_execution_bytes(WorkerBytes), Bytes is min(ParserBytes, WorkerBytes). import_worker_execution_bytes(65536). % Keep an unterminated multi-row INSERT bounded while reading it. The normal % statement scanner cannot emit anything before ';', so one 4.9 MiB INSERT % used to accumulate across read blocks and exhaust the worker before the % post-statement splitter ever had a chance to run. % The retained reverse-code list is permitted to span a couple of bounded % scanner reads so a VALUES tuple can finish, while remaining below the worker % envelope. import_partial_buffer_bytes(65536). drain_import_partial_insert(_, RevAcc0, Statements0, RevAcc, Statements) :- length(RevAcc0, BufferedBytes), import_partial_buffer_bytes(BufferLimit), BufferedBytes >= BufferLimit, !, reverse(RevAcc0, Codes), ( import_insert_values_header(Codes, Header, ValueCodes), import_values_rows_partial(ValueCodes, CompleteRows, Continuation, Delimited), CompleteRows \= [] -> import_partial_rows_to_emit(CompleteRows, Continuation, Delimited, RowsToEmit, RemainingValues), ( RowsToEmit == [] -> throw(error(resource_error(import_values_row_bytes(BufferedBytes)), context(asadb_web:drain_import_partial_insert/5, 'One VALUES row exceeds the bounded import buffer.'))) ; import_statement_batch_bytes(MaxBytes), import_chunk_insert_rows(Header, RowsToEmit, MaxBytes, Chunks), append(Header, RemainingValues, RemainingCodes), % A read block can end inside a quoted value. Recompute scanner % state from the bounded retained tail, otherwise the next block % could treat that value's closing quote as a fresh opener. scan_sql_line(RemainingCodes, none, [], [], _RemainingMode, RevAcc, []), append(Statements0, Chunks, Statements) ) ; throw(error(resource_error(import_statement_bytes(BufferedBytes)), context(asadb_web:drain_import_partial_insert/5, 'Unterminated oversized SQL must be INSERT ... VALUES rows.'))) ). drain_import_partial_insert(_, RevAcc, Statements, RevAcc, Statements). % Parse as many complete top-level VALUES tuples as are available in the % current read buffer. `Continuation` is an unfinished tuple; `Delimited` % records that a separator was consumed at the end of the buffer. import_values_rows_partial(Codes, Rows, Continuation, Delimited) :- import_skip_sql_space(Codes, First), import_values_rows_partial(First, [], Rows, Continuation, Delimited). import_values_rows_partial([], RevRows, Rows, [], false) :- !, reverse(RevRows, Rows). import_values_rows_partial(Codes, RevRows, Rows, Codes, false) :- Codes \= [40|_], !, reverse(RevRows, Rows). import_values_rows_partial(Codes, RevRows, Rows, Continuation, Delimited) :- ( import_take_values_tuple(Codes, Tuple, Rest0) -> import_skip_sql_space(Rest0, Rest), ( Rest = [] -> reverse([Tuple|RevRows], Rows), Continuation = [], Delimited = false ; Rest = [44|AfterComma] -> import_skip_sql_space(AfterComma, Next), ( Next = [] -> reverse([Tuple|RevRows], Rows), Continuation = [], Delimited = true ; import_values_rows_partial(Next, [Tuple|RevRows], Rows, Continuation, Delimited) ) ; reverse([Tuple|RevRows], Rows), Continuation = Rest, Delimited = false ) ; reverse(RevRows, Rows), Continuation = Codes, Delimited = false ). import_partial_rows_to_emit(Rows, [], Delimited, EmitRows, RemainingValues) :- !, once(append(EmitRows, [LastRow], Rows)), ( Delimited == true -> append(LastRow, [44], RemainingValues) ; RemainingValues = LastRow ). import_partial_rows_to_emit(Rows, Continuation, _, Rows, Continuation). % statement_groups_sql/2 appends a semicolon and newline to every queued % group. Count those bytes as part of the cap so the parser never receives a % payload larger than import_statement_batch_bytes/1. import_statement_sql_bytes(Codes, Bytes) :- length(Codes, Length), Bytes is Length + 2. % A generated MySQL dump can legally put thousands of VALUES tuples in one % INSERT. Such a statement cannot be admitted wholesale just because it is % syntactically complete: parse it only after slicing at top-level row tuple % boundaries. The splitter deliberately accepts only the lossless % INSERT ... VALUES (...), (...) form; another oversized statement fails with % a clear resource error instead of entering the SQL frontend unbounded. import_statement_chunks(Codes, [Codes]) :- import_statement_sql_bytes(Codes, Bytes), import_execution_batch_bytes(MaxBytes), Bytes =< MaxBytes, \+ import_insert_exceeds_row_limit(Codes), !. import_statement_chunks(Codes, Chunks) :- import_insert_values_header(Codes, Header, ValueCodes), import_values_rows(ValueCodes, Rows), import_execution_batch_bytes(MaxBytes), import_chunk_insert_rows(Header, Rows, MaxBytes, Chunks), !. import_statement_chunks(Codes, _) :- import_statement_sql_bytes(Codes, Bytes), throw(error(resource_error(import_statement_batch_bytes(Bytes)), context(asadb_web:import_statement_chunks/2, 'Oversized SQL cannot be split safely; use INSERT ... VALUES rows.'))). import_insert_exceeds_row_limit(Codes) :- import_insert_values_header(Codes, _, ValueCodes), import_values_rows(ValueCodes, Rows), import_insert_execution_row_limit(Limit), length(Rows, Count), Count > Limit. import_insert_values_header(Codes, Header, ValueCodes) :- import_find_values_keyword(Codes, none, 0, [], Header, ValueCodes), import_sql_starts_with_keyword(Header, insert). import_sql_starts_with_keyword(Codes, Keyword) :- import_skip_sql_space(Codes, Trimmed), atom_codes(Keyword, KeywordCodes), import_codes_ci_prefix(Trimmed, KeywordCodes, Rest), import_word_boundary_after(Rest). import_codes_ci_prefix(Rest, [], Rest) :- !. import_codes_ci_prefix([Code|Codes], [Expected|ExpectedCodes], Rest) :- import_ascii_lower(Code, Lower), Lower =:= Expected, import_codes_ci_prefix(Codes, ExpectedCodes, Rest). import_ascii_lower(Code, Lower) :- Code >= 65, Code =< 90, !, Lower is Code + 32. import_ascii_lower(Code, Code). import_find_values_keyword([V,A,L,U,E,S|Rest], none, 0, Rev, Header, ValueCodes) :- import_codes_ci_word([V,A,L,U,E,S], values), import_word_boundary_before(Rev), import_word_boundary_after(Rest), !, reverse([S,E,U,L,A,V|Rev], Header), ValueCodes = Rest. import_find_values_keyword([45,45|Codes], none, Depth, Rev, Header, ValueCodes) :- !, import_find_values_keyword(Codes, line_comment, Depth, [45,45|Rev], Header, ValueCodes). import_find_values_keyword([35|Codes], none, Depth, Rev, Header, ValueCodes) :- !, import_find_values_keyword(Codes, line_comment, Depth, [35|Rev], Header, ValueCodes). import_find_values_keyword([47,42|Codes], none, Depth, Rev, Header, ValueCodes) :- !, import_find_values_keyword(Codes, block_comment, Depth, [42,47|Rev], Header, ValueCodes). import_find_values_keyword([39|Codes], none, Depth, Rev, Header, ValueCodes) :- !, import_find_values_keyword(Codes, single, Depth, [39|Rev], Header, ValueCodes). import_find_values_keyword([34|Codes], none, Depth, Rev, Header, ValueCodes) :- !, import_find_values_keyword(Codes, double, Depth, [34|Rev], Header, ValueCodes). import_find_values_keyword([96|Codes], none, Depth, Rev, Header, ValueCodes) :- !, import_find_values_keyword(Codes, backtick, Depth, [96|Rev], Header, ValueCodes). import_find_values_keyword([40|Codes], none, Depth0, Rev, Header, ValueCodes) :- !, Depth is Depth0 + 1, import_find_values_keyword(Codes, none, Depth, [40|Rev], Header, ValueCodes). import_find_values_keyword([41|Codes], none, Depth0, Rev, Header, ValueCodes) :- Depth0 > 0, !, Depth is Depth0 - 1, import_find_values_keyword(Codes, none, Depth, [41|Rev], Header, ValueCodes). import_find_values_keyword([10|Codes], line_comment, Depth, Rev, Header, ValueCodes) :- !, import_find_values_keyword(Codes, none, Depth, [10|Rev], Header, ValueCodes). import_find_values_keyword([42,47|Codes], block_comment, Depth, Rev, Header, ValueCodes) :- !, import_find_values_keyword(Codes, none, Depth, [47,42|Rev], Header, ValueCodes). import_find_values_keyword([92,Code|Codes], single, Depth, Rev, Header, ValueCodes) :- !, import_find_values_keyword(Codes, single, Depth, [Code,92|Rev], Header, ValueCodes). import_find_values_keyword([92,Code|Codes], double, Depth, Rev, Header, ValueCodes) :- !, import_find_values_keyword(Codes, double, Depth, [Code,92|Rev], Header, ValueCodes). import_find_values_keyword([39|Codes], single, Depth, Rev, Header, ValueCodes) :- !, import_find_values_keyword(Codes, none, Depth, [39|Rev], Header, ValueCodes). import_find_values_keyword([34|Codes], double, Depth, Rev, Header, ValueCodes) :- !, import_find_values_keyword(Codes, none, Depth, [34|Rev], Header, ValueCodes). import_find_values_keyword([96|Codes], backtick, Depth, Rev, Header, ValueCodes) :- !, import_find_values_keyword(Codes, none, Depth, [96|Rev], Header, ValueCodes). import_find_values_keyword([Code|Codes], Mode, Depth, Rev, Header, ValueCodes) :- import_find_values_keyword(Codes, Mode, Depth, [Code|Rev], Header, ValueCodes). import_codes_ci_word(Codes, Word) :- atom_codes(Word, Expected), same_length(Codes, Expected), import_codes_ci_prefix(Codes, Expected, []). import_word_boundary_before([]). import_word_boundary_before([Code|_]) :- \+ import_identifier_code(Code). import_word_boundary_after([]). import_word_boundary_after([Code|_]) :- \+ import_identifier_code(Code). import_identifier_code(Code) :- ( Code >= 48, Code =< 57 ; Code >= 65, Code =< 90 ; Code >= 97, Code =< 122 ; Code =:= 95 ). import_skip_sql_space([Code|Codes], Rest) :- memberchk(Code, [9,10,13,32]), !, import_skip_sql_space(Codes, Rest). import_skip_sql_space(Codes, Codes). import_values_rows(Codes, Rows) :- import_skip_sql_space(Codes, First), import_values_rows(First, [], Rows). import_values_rows([], RevRows, Rows) :- !, RevRows \= [], reverse(RevRows, Rows). import_values_rows(Codes, RevRows, Rows) :- Codes = [40|_], !, import_take_values_tuple(Codes, Tuple, Rest0), import_skip_sql_space(Rest0, Rest), ( Rest = [] -> reverse([Tuple|RevRows], Rows) ; Rest = [44|AfterComma] -> import_skip_sql_space(AfterComma, Next), import_values_rows(Next, [Tuple|RevRows], Rows) ). import_take_values_tuple([40|Codes], Tuple, Rest) :- import_take_values_tuple(Codes, none, 1, [40], Tuple, Rest). import_take_values_tuple([41|Rest], none, 1, Rev, Tuple, Rest) :- !, reverse([41|Rev], Tuple). import_take_values_tuple([45,45|Codes], none, Depth, Rev, Tuple, Rest) :- !, import_take_values_tuple(Codes, line_comment, Depth, [45,45|Rev], Tuple, Rest). import_take_values_tuple([35|Codes], none, Depth, Rev, Tuple, Rest) :- !, import_take_values_tuple(Codes, line_comment, Depth, [35|Rev], Tuple, Rest). import_take_values_tuple([47,42|Codes], none, Depth, Rev, Tuple, Rest) :- !, import_take_values_tuple(Codes, block_comment, Depth, [42,47|Rev], Tuple, Rest). import_take_values_tuple([39|Codes], none, Depth, Rev, Tuple, Rest) :- !, import_take_values_tuple(Codes, single, Depth, [39|Rev], Tuple, Rest). import_take_values_tuple([34|Codes], none, Depth, Rev, Tuple, Rest) :- !, import_take_values_tuple(Codes, double, Depth, [34|Rev], Tuple, Rest). import_take_values_tuple([96|Codes], none, Depth, Rev, Tuple, Rest) :- !, import_take_values_tuple(Codes, backtick, Depth, [96|Rev], Tuple, Rest). import_take_values_tuple([40|Codes], none, Depth0, Rev, Tuple, Rest) :- !, Depth is Depth0 + 1, import_take_values_tuple(Codes, none, Depth, [40|Rev], Tuple, Rest). import_take_values_tuple([41|Codes], none, Depth0, Rev, Tuple, Rest) :- Depth0 > 1, !, Depth is Depth0 - 1, import_take_values_tuple(Codes, none, Depth, [41|Rev], Tuple, Rest). import_take_values_tuple([10|Codes], line_comment, Depth, Rev, Tuple, Rest) :- !, import_take_values_tuple(Codes, none, Depth, [10|Rev], Tuple, Rest). import_take_values_tuple([42,47|Codes], block_comment, Depth, Rev, Tuple, Rest) :- !, import_take_values_tuple(Codes, none, Depth, [47,42|Rev], Tuple, Rest). import_take_values_tuple([92,Code|Codes], single, Depth, Rev, Tuple, Rest) :- !, import_take_values_tuple(Codes, single, Depth, [Code,92|Rev], Tuple, Rest). import_take_values_tuple([92,Code|Codes], double, Depth, Rev, Tuple, Rest) :- !, import_take_values_tuple(Codes, double, Depth, [Code,92|Rev], Tuple, Rest). import_take_values_tuple([39|Codes], single, Depth, Rev, Tuple, Rest) :- !, import_take_values_tuple(Codes, none, Depth, [39|Rev], Tuple, Rest). import_take_values_tuple([34|Codes], double, Depth, Rev, Tuple, Rest) :- !, import_take_values_tuple(Codes, none, Depth, [34|Rev], Tuple, Rest). import_take_values_tuple([96|Codes], backtick, Depth, Rev, Tuple, Rest) :- !, import_take_values_tuple(Codes, none, Depth, [96|Rev], Tuple, Rest). import_take_values_tuple([Code|Codes], Mode, Depth, Rev, Tuple, Rest) :- import_take_values_tuple(Codes, Mode, Depth, [Code|Rev], Tuple, Rest). import_chunk_insert_rows(Header, Rows, MaxBytes, Chunks) :- length(Header, HeaderBytes), BaseBytes is HeaderBytes + 2, import_chunk_insert_rows(Rows, Header, MaxBytes, [], 0, BaseBytes, [], Chunks). % Keep the parsed AST and the page/checksum writer within the same bounded % unit. This is independent of the 512 KiB parser ceiling because narrow % tuples can pack thousands of rows into a much smaller SQL string. import_insert_execution_row_limit(512). import_chunk_insert_rows([], Header, _, RevRows, _, _, RevChunks, Chunks) :- !, import_finish_insert_chunk(Header, RevRows, RevChunks, FinalRevChunks), reverse(FinalRevChunks, Chunks). import_chunk_insert_rows([Row|Rows], Header, MaxBytes, [], _, BaseBytes, RevChunks, Chunks) :- !, length(Row, RowBytes), CandidateBytes is BaseBytes + RowBytes, ( CandidateBytes =< MaxBytes -> import_chunk_insert_rows(Rows, Header, MaxBytes, [Row], 1, CandidateBytes, RevChunks, Chunks) ; throw(error(resource_error(import_values_row_bytes(RowBytes)), context(asadb_web:import_chunk_insert_rows/5, 'One VALUES row exceeds the parser batch limit.'))) ). import_chunk_insert_rows([Row|Rows], Header, MaxBytes, RevRows, CurrentRows, CurrentBytes, RevChunks, Chunks) :- length(Row, RowBytes), CandidateBytes is CurrentBytes + 1 + RowBytes, import_insert_execution_row_limit(RowLimit), ( CandidateBytes =< MaxBytes, CurrentRows < RowLimit -> NextRows is CurrentRows + 1, import_chunk_insert_rows(Rows, Header, MaxBytes, [Row|RevRows], NextRows, CandidateBytes, RevChunks, Chunks) ; import_finish_insert_chunk(Header, RevRows, RevChunks, NextRevChunks), length(Header, HeaderBytes), BaseBytes is HeaderBytes + 2, import_chunk_insert_rows([Row|Rows], Header, MaxBytes, [], 0, BaseBytes, NextRevChunks, Chunks) ). import_finish_insert_chunk(Header, RevRows, RevChunks, [Statement|RevChunks]) :- reverse(RevRows, Rows), import_join_values_rows(Rows, ValueCodes), append(Header, ValueCodes, Statement). import_join_values_rows(Rows, Codes) :- import_values_row_parts(Rows, Parts), append(Parts, Codes). import_values_row_parts([], []). import_values_row_parts([Row], [Row]) :- !. import_values_row_parts([Row|Rows], [Row,[44]|Parts]) :- import_values_row_parts(Rows, Parts). import_execute_sql_batch(StopOnError, SQLCodes, Result) :- asadb_exec_sql_import_batch(SQLCodes, StopOnError, Result). release_import_batch_memory :- flag(asadb_import_batches, Batch0, Batch0 + 1), % A Reservoir worker is intentionally capped at 64 MiB. The current % SQL batch has been fully retracted before this point, so collect between % execution units rather than allowing temporary lexer/AST/page terms to % accumulate until a later emergency collection. Constraint correctness % never depends on this: its own cache has a separate fixed entry budget. garbage_collect, trim_stacks. statement_groups_sql([], []). statement_groups_sql([Codes|Groups], SQL) :- append(Codes, [59,10|Rest], SQL), statement_groups_sql(Groups, Rest). result_error_count(error(_, _), 1) :- !. result_error_count(multi(Results), Count) :- !, maplist(result_error_count, Results, Counts), sum_list(Counts, Count). result_error_count(_, 0). set_import_progress(Id, Path, Size, Bytes, Statements, Errors, Status, Message, Done) :- with_mutex(asadb_import_progress_state, ( retractall(asadb_import_progress(Id, _, _, _, _, _, _, _, _)), assertz(asadb_import_progress(Id, Path, Size, Bytes, Statements, Errors, Status, Message, Done)) )), catch(reservoir_update_progress(Id, Bytes, Statements, Errors, Status, Message, Done), _, true). record_import_failure(ImportId, DefaultPath, DefaultSize, Error) :- import_cancelled_error(ImportId, Error), !, current_import_progress(ImportId, DefaultPath, DefaultSize, Path, Size, Bytes, Statements, Errors), set_import_progress(ImportId, Path, Size, Bytes, Statements, Errors, cancelled, 'Cancellation requested; transaction rolled back.', true). record_import_failure(ImportId, DefaultPath, DefaultSize, Error) :- term_atom_safe(Error, ErrorMessage), current_import_progress(ImportId, DefaultPath, DefaultSize, Path, Size, Bytes, Statements, Errors0), Errors is Errors0 + 1, set_import_progress(ImportId, Path, Size, Bytes, Statements, Errors, error, ErrorMessage, true). import_cancelled_error(ImportId, error(reservoir_cancelled(ImportId), _)). import_cancelled_error(ImportId, reservoir_cancelled(ImportId)). current_import_progress(Id, DefaultPath, DefaultSize, Path, Size, Bytes, Statements, Errors) :- with_mutex(asadb_import_progress_state, ( asadb_import_progress(Id, Path0, Size0, Bytes0, Statements0, Errors0, _, _, _) -> Path = Path0, Size = Size0, Bytes = Bytes0, Statements = Statements0, Errors = Errors0 ; Path = DefaultPath, Size = DefaultSize, Bytes = 0, Statements = 0, Errors = 0 )). ensure_import_not_cancelled(ImportId) :- ( catch(reservoir_cancel_requested(ImportId), _, fail) -> throw(error(reservoir_cancelled(ImportId), _)) ; true ). import_progress_result(Id, table([id,path,size_bytes,bytes_read,statements,errors,status,message,done], Rows)) :- with_mutex(asadb_import_progress_state, findall([Id,Path,Size,Bytes,Statements,Errors,Status,Message,Done], asadb_import_progress(Id, Path, Size, Bytes, Statements, Errors, Status, Message, Done), Found)), ( Found = [] -> Rows = [[Id,'',0,0,0,0,waiting,'',false]] ; Rows = Found ). merge_import_last(PrevStatus, PrevMessage, none, _, PrevStatus, PrevMessage) :- !. merge_import_last(_, _, Status, Message, Status, Message). scan_sql_line([], line_comment, RevAcc, RevStatements, none, RevAcc, Statements) :- !, reverse(RevStatements, Statements). scan_sql_line([], Mode, RevAcc, RevStatements, Mode, RevAcc, Statements) :- !, reverse(RevStatements, Statements). scan_sql_line([59|Cs], none, RevAcc, RevStatements, Mode, RevOut, Statements) :- !, reverse(RevAcc, Codes), trim_sql_codes(Codes, Trimmed), scan_sql_line(Cs, none, [], [Trimmed|RevStatements], Mode, RevOut, Statements). scan_sql_line([45,45|Cs], none, RevAcc, RevStatements, Mode, RevOut, Statements) :- !, scan_sql_line(Cs, line_comment, [45,45|RevAcc], RevStatements, Mode, RevOut, Statements). scan_sql_line([35|Cs], none, RevAcc, RevStatements, Mode, RevOut, Statements) :- !, scan_sql_line(Cs, line_comment, [35|RevAcc], RevStatements, Mode, RevOut, Statements). scan_sql_line([47,42|Cs], none, RevAcc, RevStatements, Mode, RevOut, Statements) :- !, scan_sql_line(Cs, block_comment, [42,47|RevAcc], RevStatements, Mode, RevOut, Statements). scan_sql_line([39|Cs], none, RevAcc, RevStatements, Mode, RevOut, Statements) :- !, scan_sql_line(Cs, single, [39|RevAcc], RevStatements, Mode, RevOut, Statements). scan_sql_line([34|Cs], none, RevAcc, RevStatements, Mode, RevOut, Statements) :- !, scan_sql_line(Cs, double, [34|RevAcc], RevStatements, Mode, RevOut, Statements). scan_sql_line([96|Cs], none, RevAcc, RevStatements, Mode, RevOut, Statements) :- !, scan_sql_line(Cs, backtick, [96|RevAcc], RevStatements, Mode, RevOut, Statements). scan_sql_line([10|Cs], line_comment, RevAcc, RevStatements, Mode, RevOut, Statements) :- !, scan_sql_line(Cs, none, [10|RevAcc], RevStatements, Mode, RevOut, Statements). scan_sql_line([42,47|Cs], block_comment, RevAcc, RevStatements, Mode, RevOut, Statements) :- !, scan_sql_line(Cs, none, [47,42|RevAcc], RevStatements, Mode, RevOut, Statements). scan_sql_line([92,C|Cs], single, RevAcc, RevStatements, Mode, RevOut, Statements) :- !, scan_sql_line(Cs, single, [C,92|RevAcc], RevStatements, Mode, RevOut, Statements). scan_sql_line([92,C|Cs], double, RevAcc, RevStatements, Mode, RevOut, Statements) :- !, scan_sql_line(Cs, double, [C,92|RevAcc], RevStatements, Mode, RevOut, Statements). scan_sql_line([39|Cs], single, RevAcc, RevStatements, Mode, RevOut, Statements) :- !, scan_sql_line(Cs, none, [39|RevAcc], RevStatements, Mode, RevOut, Statements). scan_sql_line([34|Cs], double, RevAcc, RevStatements, Mode, RevOut, Statements) :- !, scan_sql_line(Cs, none, [34|RevAcc], RevStatements, Mode, RevOut, Statements). scan_sql_line([96|Cs], backtick, RevAcc, RevStatements, Mode, RevOut, Statements) :- !, scan_sql_line(Cs, none, [96|RevAcc], RevStatements, Mode, RevOut, Statements). scan_sql_line([C|Cs], Mode0, RevAcc, RevStatements, Mode, RevOut, Statements) :- scan_sql_line(Cs, Mode0, [C|RevAcc], RevStatements, Mode, RevOut, Statements). execute_statement_codes([], _, Statements, Errors, Statements, Errors, none, '', false). execute_statement_codes([Codes|Rest], StopOnError, Statements0, Errors0, Statements, Errors, LastStatus, LastMessage, Stop) :- ( Codes == [] -> execute_statement_codes(Rest, StopOnError, Statements0, Errors0, Statements, Errors, LastStatus, LastMessage, Stop) ; atom_codes(SQL, Codes), asadb_exec_sql(SQL, Result), result_status_message(Result, Status, Message), Statements1 is Statements0 + 1, ( result_has_error(Result) -> Errors1 is Errors0 + 1 ; Errors1 = Errors0 ), ( StopOnError == true, result_has_error(Result) -> Statements = Statements1, Errors = Errors1, LastStatus = Status, LastMessage = Message, Stop = true ; execute_statement_codes(Rest, StopOnError, Statements1, Errors1, Statements, Errors, LastStatus0, LastMessage0, Stop), ( LastStatus0 == none -> LastStatus = Status, LastMessage = Message ; LastStatus = LastStatus0, LastMessage = LastMessage0 ) ) ). result_has_error(error(_, _)) :- !. result_has_error(multi(Results)) :- !, member(Result, Results), result_has_error(Result). result_has_error(_):- fail. result_status_message(ok(Message), ok, Message) :- !. result_status_message(error(_, Message), error, Message) :- !. result_status_message(table(_, Rows), table, Message) :- !, length(Rows, Count), format(atom(Message), '~w row(s)', [Count]). result_status_message(multi(Results), Status, Message) :- !, ( first_result_error_status_message(Results, Status, Message) -> true ; last_result_status_message(Results, Status, Message) ). result_status_message(Result, ok, Atom) :- term_to_atom(Result, Atom). last_result_status_message([], none, ''). last_result_status_message([Result], Status, Message) :- !, result_status_message(Result, Status, Message). last_result_status_message([_|Rest], Status, Message) :- last_result_status_message(Rest, Status, Message). first_result_error_status_message([Result|_], Status, Message) :- result_has_error(Result), !, result_status_message(Result, Status, Message). first_result_error_status_message([_|Results], Status, Message) :- first_result_error_status_message(Results, Status, Message). trim_sql_codes(Codes, Trimmed) :- drop_sql_ws(Codes, Left), reverse(Left, Rev), drop_sql_ws(Rev, RevTrimmed), reverse(RevTrimmed, Trimmed). drop_sql_ws([C|Cs], Rest) :- member(C, [9,10,13,32]), !, drop_sql_ws(Cs, Rest). drop_sql_ws(Codes, Codes). catalog_row(CurrentDb, DBs, [CurrentDb, DBName, database, DBName, 0, [], [], '']) :- member(DBTerm, DBs), db_catalog_parts(DBTerm, DBName, _, _), visible_catalog_db(DBName). catalog_row(CurrentDb, DBs, [CurrentDb, DBName, table, TableName, RowCount, Columns, Indexes, '']) :- member(DBTerm, DBs), db_catalog_parts(DBTerm, DBName, Tables, _), visible_catalog_db(DBName), member(TableTerm, Tables), table_catalog_parts(TableTerm, TableName, Columns, RowStorage, Indexes), catalog_row_count(RowStorage, RowCount). catalog_row(CurrentDb, DBs, [CurrentDb, DBName, view, ViewName, 0, [], [], Query]) :- member(DBTerm, DBs), db_catalog_parts(DBTerm, DBName, _, Views), visible_catalog_db(DBName), member(ViewTerm, Views), view_catalog_parts(ViewTerm, ViewName, Query). visible_catalog_db(Name) :- \+ sub_atom(Name, 0, 2, _, '__'). db_catalog_parts(db(Name, Tables, Views, _, _, _), Name, Tables, Views) :- !. db_catalog_parts(db(Name, Tables), Name, Tables, []). table_catalog_parts(table(Name, Columns, Rows, Indexes), Name, Columns, Rows, Indexes) :- !. table_catalog_parts(table(Name, Columns, Rows), Name, Columns, Rows, []). catalog_row_count(paged_rows(_, Count, _), Count) :- !. catalog_row_count(Rows, Count) :- length(Rows, Count). view_catalog_parts(view(Name, Query, _), Name, Query) :- !. view_catalog_parts(view(Name, Query), Name, Query) :- !. view_catalog_parts(Name, Name, ''). read_sql_body(Request, SQL) :- http_read_data(Request, Data, []), member(sql=SQL, Data), atom_length(SQL, Len), Len =< 250000. request_stream_payload(Request, In, Size) :- member(input(RawIn), Request), request_content_size(Request, Size), asadb_config_get(reservoir_max_spool_bytes, MaxSpoolBytes), Size > 0, Size =< MaxSpoolBytes, stream_range_open(RawIn, In, [size(Size)]), set_stream(In, encoding(octet)). request_content_size(Request, Size) :- member(content_length(Size0), Request), content_size_number(Size0, Number), !, Size is max(0, Number). request_content_size(_, 0). content_size_number(Value, Value) :- number(Value), !. content_size_number(Value, Number) :- atom(Value), atom_number(Value, Number). request_import_id(Request, Id) :- request_header(x_asadb_import_id, Request, Raw), atom(Raw), Raw \= '', !, Id = Raw. request_import_id(_, Id) :- random_between(100000000, 999999999, N), format(atom(Id), 'stream-~w', [N]). request_stop_on_error(Request, true) :- request_header(x_asadb_stop_on_error, Request, Value), member(Value, [true, 'true', yes, 'yes', '1', 1]), !. request_stop_on_error(_, false). request_idempotency_key(Request, Key) :- request_header(x_asadb_idempotency_key, Request, Raw), atom(Raw), Raw \= '', !, Key = Raw. request_idempotency_key(Request, Key) :- request_header(x_asadb_import_id, Request, Raw), atom(Raw), Raw \= '', !, Key = Raw. request_idempotency_key(_, ''). request_job_label(Request, Label) :- request_header(x_asadb_job_label, Request, Raw), atom(Raw), Raw \= '', !, Label = Raw. request_job_label(_, 'SQL command'). request_reservoir_metadata(Request, Metadata) :- request_reservoir_database(Request, Database), request_header(x_asadb_import_format, Request, Format), atom(Format), Format \== '', !, request_header_default(x_asadb_import_name, Request, 'import.sql', SourceName), request_header_default(x_asadb_import_table, Request, '', Target), request_header_default(x_asadb_import_mode, Request, replace, Mode), Metadata = _{ kind:interchange, format:Format, source_name:SourceName, target_table:Target, mode:Mode, logical_database:Database }. request_reservoir_metadata(Request, Metadata) :- request_reservoir_database(Request, Database), Metadata = _{logical_database:Database}. request_reservoir_database(Request, Database) :- request_header_default(x_asadb_logical_database, Request, '', Header), ( Header == '' -> asadb_current_database(Database) ; Database = Header ). request_header_default(Name, Request, Default, Value) :- ( request_header(Name, Request, Found), atom(Found), Found \== '' -> Value = Found ; Value = Default ). form_idempotency_key(Data, Key) :- member(idempotency_key=Raw, Data), atom(Raw), Raw \= '', !, Key = Raw. form_idempotency_key(Data, Key) :- import_id(Data, Key). reservoir_http_job_id(Request, JobId) :- http_parameters(Request, [id(JobId, [atom])]). reservoir_http_result_parameters(Request, JobId, Offset, Limit) :- asadb_config_get(reservoir_result_page_rows, DefaultLimit), http_parameters(Request, [ id(JobId, [atom]), offset(Offset0, [integer, default(0)]), limit(Limit0, [integer, default(DefaultLimit)]) ]), Offset is max(0, Offset0), Limit is max(1, min(DefaultLimit, Limit0)). reservoir_api_error(error(resource_error(reservoir_job_slots), _)) :- !, json_error('429 Too Many Requests', 'Reservoir queue is full; retry after an active job finishes.'). reservoir_api_error(error(resource_error(reservoir_spool_capacity), _)) :- !, json_error('507 Insufficient Storage', 'Reservoir spool capacity would be exceeded.'). reservoir_api_error(error(existence_error(reservoir_job, _), _)) :- !, json_error('404 Not Found', 'Reservoir job was not found or has expired.'). reservoir_api_error(Error) :- term_atom_safe(Error, Message), json_error('400 Bad Request', Message). authorized_api(Request) :- request_token(Request, Token), asadb_panel_token(Token), !. request_token(Request, Token) :- request_header(x_asadb_token, Request, Token). request_token(Request, Token) :- request_header(cookie, Request, Cookie), cookie_token(Cookie, Token). request_header(Name, Request, Value) :- member(Header, Request), Header =.. [Name, Value]. cookie_token(Cookie, Token) :- is_list(Cookie), !, member(asadb_token=Token, Cookie). cookie_token(Cookie, Token) :- atomic_list_concat(Parts, ';', Cookie), member(Part0, Parts), normalize_space(atom(Part), Part0), sub_atom(Part, 0, _, _, 'asadb_token='), sub_atom(Part, 12, _, 0, Token). security_headers :- format('X-Content-Type-Options: nosniff~n'), format('Referrer-Policy: no-referrer~n'), format('X-Frame-Options: DENY~n'), format('Cross-Origin-Resource-Policy: same-origin~n'), format('Content-Security-Policy: default-src ''self''; script-src ''self''; style-src ''self''; img-src ''self'' data:; media-src ''self''; connect-src ''self''; base-uri ''none''; frame-ancestors ''none''~n'), format('Cache-Control: no-store~n'). json_response(JSON) :- security_headers, format('Content-type: application/json~n~n'), format('~w', [JSON]). json_dict_response(Dict) :- security_headers, format('Content-type: application/json~n~n'), json_write_dict(current_output, Dict, [width(0)]). json_dict_response(Status, Dict) :- format('Status: ~w~n', [Status]), json_dict_response(Dict). json_error(Status, Message) :- format('Status: ~w~n', [Status]), security_headers, format('Content-type: application/json~n~n'), json_write_dict(current_output, _{status:error,message:Message}, [width(0)]).