diff --git a/docs/epics.md b/docs/epics.md index cff4baa..4adbaf6 100644 --- a/docs/epics.md +++ b/docs/epics.md @@ -26,6 +26,7 @@ Planned epics not yet started: - **Query Federation**: Add support for federated queries across datacenters - **Tag-based Specifiers**: Allow tag-based query federation +- **Async Query Optimization**: True async I/O implementation for query commands to improve performance (currently uses thread pools) - **Dataset Cleaning**: Tools for cleaning/validating datasets - **Overall Clean-up**: Code restructuring and cleanup diff --git a/docs/source/commands.md b/docs/source/commands.md index ef0ab15..b25cb8f 100644 --- a/docs/source/commands.md +++ b/docs/source/commands.md @@ -1,501 +1,64 @@ # OWIlix Commands and Usage -The OWIlix CLI includes custom commands that can be executed via the CLI. -Commands are utilising core functions and add a python rich ui based console UI on top of it. -Regular python functions can be registered as commands, where the parameters are translated to CLI parameters and the function name becomes the command name. -The main components are shown in the figure below. +The OWIlix CLI provides a suite of commands for managing datasets, interacting with remote repositories, and performing data analysis. -![](_static/owilix-commands.png) +## Command Groups -## Command-Overview +| Group | Description | Documentation | +|-------|-------------|---------------| +| **[local](details/local.md)** | Manage datasets in your local repository. | [Read More](details/local.md) | +| **[remote](details/remote.md)** | Interact with remote data centers (pull, push, list). | [Read More](details/remote.md) | +| **[query](details/query.md)** | SQL-based analysis and filtering (slice, stats, less). | [Read More](details/query.md) | +| **[config](details/config.md)** | View and modify CLI configuration. | [Read More](details/config.md) | +| **admin** | Administrative tasks (logs, system stats). | *See below* | -| Group | Command | Description | -|-------|---------|-------------| -| **local** | `ls` | List local datasets | -| | `rm` | Remove local datasets | -| | `init` | Initialize new dataset | -| | `insert` | Insert files into dataset | -| | `free` | Free dataset space | -| | `export` | Export dataset to archive | -| | `analyze` | Analyze JSONL files | -| **remote** | `ls` | List remote datasets | -| | `pull` | Download datasets | -| | `push` | Upload datasets | -| | `diff` | Compare local/remote | -| | `doctor` | Check connections | -| | `logout` | Revoke tokens | -| **query** | `less` | Interactive SQL browser | -| | `sites` | URL-based filtering | -| | `slice` | Create new dataset from query | -| | `stream` | Stream results to consumer | -| | `warc` | Extract WARC file locations | -| | `aggregate` | SQL aggregation | -| | `analyze` | Job log analysis | -| **config** | `version` | Show config version | -| | `list` | Show configuration | -| | `get` | Get config values | -| | `set` | Set config values | -| **admin** | `logs` | View error logs | -| | `stats` | Usage statistics | -| | `check` | Validate dataset metadata | -| **batch** | `run` | Execute batch jobs from YAML | +## Detailed Documentation -## Usage +- **[Local Commands](details/local.md)**: `ls`, `rm`, `analyze`, `init` +- **[Remote Commands](details/remote.md)**: `ls`, `pull`, `push`, `diff`, `doctor`, `logout` +- **[Query Commands](details/query.md)**: + - **[Slice](details/query_slice.md)**: Create new datasets from queries. + - **[WARC](details/warc.md)**: WARC file extraction. + - **[Graphs](details/graphs.md)**: Web graph analysis. + - `less` (interactive data browser), `stats`, `sites`, `aggregate` + - Supports `--async` mode for concurrent query execution +- **[Configuration](details/config.md)**: `list`, `get`, `set` -### Defaults -- **Path**: `~/.owi` (modifiable via `OWS_OWI_PATH` environment variable or using the `--target` option) -- **File Names**: `{internalid}.tar.gz` for datasets and `{internalid}.json` for metadata -- **Specifier Format**: `{datacenter|all}:{YYYY-MM-DD|latest}#{days}/{key=value;key=value}` +## Global Options -#### Specifier Components -1. **Datacenter**: Specific data center or "all" (default: "all"). Current datacenters are `lrz`,`it4i`, `csc` -2. **Date**: In `YYYY-MM-DD` or "latest" (default: current day) -3. **Days**: Days in the past to retrieve (optional) -4. **Key=Value**: Additional filters (optional) - - **Exact match**: `key=value` - matches exact value - - **Regex match**: `key*=pattern` - uses regex pattern matching - -##### Specifier Examples - -```bash -# Exact match -all:latest/collectionName=cefal - -# Regex match - all collections -all:latest/collectionName*=.* - -# Regex match - collections starting with "gpu" -lrz:latest/collectionName*=gpu.* - -# Regex match - titles containing "legal" -it4i:latest/title*=.*legal.* - -# Combined filters -all:latest/access=public/collectionName*=main.* -``` - -### Global Flags - -```bash -owi [OPTIONS] COMMAND [ARGS]... -``` +The following flags can be used with almost any command: | Flag | Short | Description | |------|-------|-------------| -| `--verbose` | `-v` | Enable verbose output | -| `--loglevel` | `-l` | Log level: DEBUG, INFO, WARNING, ERROR | -| `--yes` | `-y` | Auto-confirm all prompts | -| `--config` | `-c` | Configuration file path | -| `--format` | `-f` | Output format: table, json, jsonl, json+, jsonl+ | -| `--output` | `-o` | Write output to file | -| `--no-display` | | Suppress console output | -| `--target` | `-t` | Target directory for data | -| `--no-progress` | `-N` | Suppress progress bars | - -### Commands and Examples - -Note that `owi` and `owilix` commands are installed. - -`owilix` defines different command groups: local, remote, query, config, admin - -#### Login and Logout - -`owilix` requires you to login using B2ACCESS when working with remote repositories, i.e. you have to follow the link and complete the device login. Afterwards it retains a token with a specific refresh timeout. - -To revoke that token from the machine, you must use: -```bash -owi remote logout -``` - -#### Local Commands - -Local commands `owi local` are used to manage local datasets. - -##### Listing local datasets - -```bash -owi local ls all -owi local ls all:latest --display wide -owi local ls all:latest#14/access=public --sort totalSize --reverse -``` - -Display formats: `short` (default), `wide` (panels), `full` (paged YAML), `markdown` - -Output columns in short format: ID, Title, Date, Collection, DataCenter, Zone, ObjectCount, Size, Access - -##### Removing datasets - -```bash -owi local rm "all/id=abc123..." -owi local rm old_collection:2023-01 -y # Skip confirmation -``` - -##### Analyzing JSONL files - -```bash -owi local analyze data.jsonl -owi local analyze data.jsonl.gz --topk 50 -``` - -> **Note**: `local insert` and `local export` are available in the legacy CLI. - -#### Query Commands - -Query commands allow SQL-based interaction with datasets. They use `-L/--local` and `-R/--remote` options to specify targets. - -##### Output Formats - -Query commands support the global `--format` flag with the following options: - -| Format | Description | -|--------|-------------| -| `table` | Interactive table with pagination (default) | -| `json` | JSON output, one object per result | -| `jsonl` | JSON Lines format (newline-delimited JSON) | -| `json+` | JSON with progress items (name: "progress-item") | -| `jsonl+` | JSONL with progress items for streaming | - -Progress items in `+` formats have `"name": "progress-item"` to distinguish from data. - -##### query less - Interactive Browser - -```bash -# Basic usage (table output, interactive pagination) -owi query less -R lrz:latest -owi query less -R lrz:latest --select "url,title" --limit 100 - -# Auto-confirm pagination with --yes -owi --yes query less -R lrz:latest --limit 50 - -# JSON output (no pagination, direct streaming) -owi --format jsonl query less -R lrz:latest --limit 100 - -# JSON with progress for pipeline processing -owi --format jsonl+ query less -R lrz:latest | jq 'select(.name != "progress-item")' - -# Output to file -owi --format jsonl -o results.jsonl query less -R lrz:latest --select "url,title" -``` - -##### query sites - URL-based Filtering - -```bash -# Filter by domains in file -owi query sites urls.txt -R lexis:latest - -# JSON output with specific columns -owi --format jsonl query sites sites.csv -R all/id=0350fecc-e58b-11f0-a8c9-8ebf6bb2cab9 --select "url,title" - -# Output to file with progress -owi --format jsonl+ -o output.jsonl query sites sites.csv -R lexis:latest -``` - -##### query stats - Dataset Statistics - -Calculates statistics over datasets using DuckDB queries: record counts, URL distributions, date ranges, MIME types. - -```bash -# Table output with progress -owi query stats -R lexis:latest/collectionName=cefal - -# JSON output -owi --format json query stats -R all/id=0350fecc-e58b-11f0-a8c9-8ebf6bb2cab9 -o stats.json - -# Adjust top-k items -owi query stats -R lexis:latest --topk 50 -``` - -##### query aggregate - SQL Aggregation - -Run SQL aggregation queries with GROUP BY support. - -```bash -# Count URLs by suffix -owi query aggregate -R lexis:latest --select "url_suffix,count(*) as c" --groupby url_suffix - -# With custom final query -owi query aggregate -R lexis:latest --select "url_domain,count(*) as cnt" --groupby url_domain \ - --aggregate "SELECT * FROM aggregates ORDER BY cnt DESC LIMIT 10" - -# JSON output -owi --format jsonl query aggregate -R lexis:latest --select "mime_type,count(*) as c" --groupby mime_type -``` - - -#### Remote Commands - -Remote commands operate on the datacenters specified. - -##### Listing remote datasets - -```bash -owi remote ls all -owi remote ls lexis:latest --display wide -owi remote ls all/id=0350fecc-e58b-11f0-a8c9-8ebf6bb2cab9 - -# List files in a dataset (new --files flag) -owi remote ls lexis:latest/collectionName=cefal --files "**/*.parquet" -``` - -##### Listing files in datasets - -Use `--files` (`-f`) to list files matching a glob pattern: - -```bash -# List parquet files in a local dataset -owi local ls all/id=abc123 --files "**/*.parquet" - -# List parquet files in a remote dataset -owi remote ls lexis:latest/collectionName=cefal --files "**/*.parquet" -``` - -##### Pulling datasets - -```bash -owi remote pull lexis:latest/access=public -owi remote pull all/id=abc123... --files "**/*.parquet" -owi remote pull it4i:latest --language eng --threads 4 -``` - - -##### Pushing datasets - -```bash -owi remote push all/access=project --datacenter it4i -owi remote push all/id=abc123... --yes -``` - -##### Comparing local/remote (diff) - -```bash -owi remote diff all/access=project -owi remote diff all --files "**/*" # File-level diff -``` - -##### Connection check - -```bash -owi remote doctor -``` - -##### Logout - -```bash -owi remote logout -``` - -### Working with the Web-Graph - -Owilix supports the creation of a web-graph statistics and reports stored per dataset and can potentially aggregate over datasets. - -Graph data is stored in `.stats/graph/` containining a duckdb database with links and nodes (for hosts and domains) and a `graph_report.html` file if creatd. - - -- **Creating the Web Graph** The following command creates the web-graph. With update_mode you control if the web-graph is created or updated at the remote site or (if avoided) created locally -```bash -owilix query create_graphs --remote it4i/id=76381d62-561b-11f0-ab83-528c047b29ff update_mode=create -``` -Note that storing data along side the dataset requires project rights! - -- **Creating the Web Graph report** Create a HTML report with a table on node measures and a heatmap visualisation of teh topk nodes according to page rank. - -```bash -owilix query graph_reports --remote it4i/id=76381d62-561b-11f0-ab83-528c047b29ff mode=create -``` - -```bash -owilix query graph_reports --remote it4i/id=76381d62-561b-11f0-ab83-528c047b29ff mode=view -``` - -Note that report creation requires that `.stats/graph` exists. - - - -### Workflows for creating datasets - -Datasets can be created in two steps: - -1. Inserting dataset into the local repository: - - ```bash - owilix local insert /Users/username/mydata//2023-12-03 access=project collectionName="main" move=False - ``` - - After inserting, you can still add some more metadata to the local json file, or you can add additional files in the repository. - - Metadata can be overwritten during insert. - ```bash - owilix local insert /Users/mgrani/owseudata/migration/it4i/2023-12-03 access=project collectionName="main" move=False owner="OpenWebSearch.eu Consortium" creator="OpenWebSearch.eu Consortium" publisher="OpenWebSearch.eu Consortium" - ``` - -2. Pushing the data to the server. Note that this changes the internal id. - ```bash - owilix remote push it4i:latest/access=project;internalId=6ebaf89d-adf2-4680-ab67-82b8eaf999fa - ``` - Pushing includes a selection of local datasets and the specification of the remote properties, particularly the data center. - Note that in this case the datacenter in the dataset specifier is interpreted as being the remote datacenter. - This can be also set explicitly by using `dataCenter=`. Further remarks: - - you can specify a file glob using `files=glob` and only those files will be considered for syncing - - sync does not push files already at the server. This behaviour can be changed using `overwrite=True` - - Be careful when pushing data to not pollute datasets server side. - - -### Query Commands - -Query commands allow SQL-based analysis of parquet datasets using DuckDB. - -##### Interactive Browsing (less) - -```bash -owi query less main:latest "" --select "url,title" --limit 100 -owi query less "" lrz:latest --where "url_suffix='at'" --json -owi query less main:latest lrz:latest --files "**/*.parquet" --page-size 20 -``` - -##### Site Filtering (sites) - -```bash -owi query sites "" lrz:latest urls.txt --select "url,title" --json -owi query sites main:latest "" domains.txt --output results.jsonl -``` - -##### Aggregation - -```bash -owi query aggregate "" lrz:latest --select "url_suffix,count(*) as c" --groupby url_suffix -owi query aggregate main:latest "" --aggregate "SELECT * FROM aggregates ORDER BY c DESC LIMIT 20" -``` - -##### Job Analysis - -```bash -owi query analyze my_job_name --details --error-limit 50 -``` - -### Configuration Commands - -```bash -owi config list # Show all config -owi config get user.name user.email # Get specific values -owi config set user.name="John" user.email="j@x.com" # Set values -``` - -### Admin Commands - -```bash -owi admin logs error --lines 100 # View error logs -owi admin stats # Usage statistics -``` - -### Plugin Commands - -```bash -owi plugin list # Show available plugins -owi plugin run owilix.plugins.push.opensearch.OpenSearchIndexer index my_index -owi plugin run owilix.plugins.ngram.BloomNGrams create "all:latest/collectionName=main" -``` - -### Filtering for a list of websites - -`owilix` can be used to filter for list of websites given as text file. - -```bash -echo "uni-passau.de\nopensearchfoundation.org\n"> ~/tmp/list-of-urls.txt -query sites --remote all "files=**/language=eng/*.parquet" "select=url,title,curlielabels_en" urls_file=~/tmp/list-of-urls.txt -``` - -Note that since this is a query, it can take some time. - -Another example: - -```bash -owilix query sites --remote up:2025-01-15#15/collectionName=main "select=url,title,curlielabels_en,warc_date,warc_file" urls_file=list.csv as_json=True json_file=list-result.json -``` - -You can get some statistics of the resulting json using: - -```bash -owilix local analyze_jsonl list-result.json topk=100 -``` - -### Query specific configurations - -Querying has often longer timeouts. to account for that, we are setting the irods connection timeout to None. You can use the "OWILIX_IRODS_CONNECTION_TIMEOUT" to a different value. See code - -```python -fs.session.connection_timeout = os.getenv("OWILIX_IRODS_CONNECTION_TIMEOUT", None) -``` - - -## Configuration - -`owilix` relies on a yaml based configuration file. -The default directory for the yaml configuration file is `~/.owi/owilix.cfg`. -However, the `owi` directory and consequently the used config file can be changed using the `--target` option. -For changing the configuration, you can edit the file directly or also use the owilix config command: - -1. List the config and its location: `owilix config list` -2. Set a (list of) config entry: `owilix config set repositories.config.local.options.path=/home/joe/.owi/{access}` -3. get a list of entries: `owilix config get repositories.config server-config` - -The local configuration is patched upon every start with a server call to `openwebindex.eu` containing the newest configuration for remote repositories. - -You can avoid this patching by using the flag `--no-remote-config`. - -**Examples** - -Setting your identity: - -```bash -owilix config set "user.name=Michael Granitzer" "user.email=michael.granitzer@uni-passau.de" "user.organistion=University of Passau" -owilix config get user.name user.email user.organistion -``` - -Setting the default path for the local repository: - -```bash -owilix config set repositories.config.local.options.path=/home/joe/.owi/{access} -``` - - - -## OWILIX Command Modules - -```{automodule} owilix.cmd -:members: -:undoc-members: -:private-members: -:special-members: __init__ -:show-inheritance: -``` - -### Local OWILIX Commands - - -```{automodule} owilix.cmd.local.LocalCommands -:members: -:undoc-members: -:private-members: -:special-members: __init__ -:show-inheritance: -``` - -```{autofunction} owilix.cmd.local.details -``` - -```{autofunction} owilix.cmd.local.ls -``` - -```{autofunction} owilix.cmd.local.insert -``` - - -### OWILIX BaseCommands - -```{automodule} owilix.cmd.base.BaseCommand -:members: __init__ -:undoc-members: -:private-members: -:special-members: __init__ -:show-inheritance: -``` \ No newline at end of file +| `--verbose` | `-v` | Enable verbose output (debug logs). | +| `--loglevel` | `-l` | Set log level: DEBUG, INFO, WARNING, ERROR. | +| `--yes` | `-y` | Auto-confirm interactive prompts. | +| `--config` | `-c` | Specify a custom configuration file path. | +| `--format` | `-f` | Output format: `table`, `json`, `jsonl`. | +| `--output` | `-o` | Write output to a file. | +| `--no-display` | | Suppress console output (useful for scripting). | +| `--target` | `-t` | Override the target directory for data. | +| `--no-progress` | `-N` | Suppress progress bars. | + +## Defaults & Specifiers + +### Paths +- Default Path: `~/.owi` +- Can be modified via environment variable `OWS_OWI_PATH` or `owi config`. + +### Specifier Format +Access datasets using the following specifier format: +`{datacenter|all}:{YYYY-MM-DD|latest}#{days}/{key=value;key=value}` + +Examples: +- `all:latest`: Latest datasets from all datacenters. +- `lexis:2025-01-01`: Datasets from LEXIS on a specific date. +- `all/collectionName=main`: Filter by collection name. +- `all/access=public`: Filter by access level. + +## Admin Commands + +Admin commands are used for troubleshooting and system monitoring. + +- `owi admin logs [level]`: View logs. +- `owi admin stats`: View system statistics. \ No newline at end of file diff --git a/docs/source/conf.py b/docs/source/conf.py index 10aef10..ca62ffa 100644 --- a/docs/source/conf.py +++ b/docs/source/conf.py @@ -46,6 +46,6 @@ html_static_path = ['_static'] # Optional: Theme-specific options (for dropdown) html_theme_options = { - 'version_dropdown': True, # Required for version dropdown in html + # 'version_dropdown': True, # Removed as unsupported by sphinxdoc theme # Add other theme-specific options here } \ No newline at end of file diff --git a/docs/source/configuration.md b/docs/source/configuration.md index 29d60d2..a836867 100644 --- a/docs/source/configuration.md +++ b/docs/source/configuration.md @@ -248,7 +248,7 @@ To skip merging for a session, construct with `no_remote=True`. To permanently d The configuration is managed by the `OWIlixManager` and `OWIlixConfig` classes in the module `owilix.manager`. -```{automodule} owilix.manager +```{automodule} owilix.core.manager :members: :undoc-members: :private-members: diff --git a/docs/source/details/config.md b/docs/source/details/config.md new file mode 100644 index 0000000..4676317 --- /dev/null +++ b/docs/source/details/config.md @@ -0,0 +1,68 @@ +# Configuration Commands + +Configuration commands (`owi config`) allow you to view and modify the OWIlix configuration. + +## `config list` + +Show the full current configuration, merging system defaults, user config file, and remote configuration (if enabled). + +### Usage + +```bash +owi config list +``` + +--- + +## `config get` + +Get the value of one or more configuration keys. + +### Usage + +```bash +owi config get KEY [KEY...] +``` + +### Examples + +```bash +owi config get user.name +owi config get repositories.config.local.options.path +``` + +--- + +## `config set` + +Set the value of configuration keys. These are saved to your local config file (usually `~/.owi/owilix.cfg`). + +### Usage + +```bash +owi config set KEY=VALUE [KEY=VALUE...] +``` + +### Examples + +**Set user identity:** +```bash +owi config set user.name="Jane Doe" user.email="jane@example.com" +``` + +**Change local repository path:** +```bash +owi config set repositories.config.local.options.path="/data/owi/{access}" +``` + +--- + +## `config version` + +Show the version of the configuration schema and the CLI. + +### Usage + +```bash +owi config version +``` diff --git a/docs/source/details/graphs.md b/docs/source/details/graphs.md index 98b76ac..c1c3433 100644 --- a/docs/source/details/graphs.md +++ b/docs/source/details/graphs.md @@ -1,7 +1,9 @@ # Graphs ```{eval-rst} -.. automodule:: owilix.cmd.subcmds.query_graph +.. automodule:: owilix.core.tasks.query_graphs + :members: + :undoc-members: :show-inheritance: ``` diff --git a/docs/source/details/local.md b/docs/source/details/local.md new file mode 100644 index 0000000..8e16163 --- /dev/null +++ b/docs/source/details/local.md @@ -0,0 +1,102 @@ +# Local Commands + +Local commands (`owi local`) are used to manage datasets stored in your local repository. The local repository path defaults to `~/.owi` but can be configured. + +## `local ls` + +List datasets available in the local repository. + +### Usage + +```bash +owi local ls [SPECIFIER] [OPTIONS] +``` + +### Options + +- `--files GLOB`: List files within the dataset(s) matching the glob pattern (e.g., `**/*.parquet`). +- `--file-details`: When used with `--files`, shows file sizes (slower as it fetches metadata for each file). +- `--sort FIELD`: Sort output by a metadata field (e.g., `totalSize`, `date`). +- `--reverse`: Reverse the sort order. +- `--display FORMAT`: Output format for the table (`short`, `mid`, `wide`, `full`, `markdown`). default: `short`. +- `--detailed / --no-detailed`: Show detailed file counts and sizes (checks internal file structure). + +### Examples + +**List all local datasets containing "web" in the title:** +```bash +owi local ls "all/title*=.*Web.*" +``` + +**List datasets sorted by size (descending):** +```bash +owi local ls all --sort totalSize --reverse +``` + +**List all parquet files in a specific dataset:** +```bash +owi local ls all/id=123... --files "**/*.parquet" +``` + +**List files with sizes (slower):** +```bash +owi local ls all/id=123... --files "**/*.parquet" --file-details +``` + +**JSON output for machine processing:** +```bash +owi --format json local ls all +``` + +--- + +## `local rm` + +Permanently remove datasets from the local repository. + +### Usage + +```bash +owi local rm [SPECIFIER] [OPTIONS] +``` + +### Options + +- `--yes` (`-y`): Skip confirmation prompt. + +### Examples + +**Remove a specific dataset by ID:** +```bash +owi local rm "all/id=a739b3cf-11c4-4c81-ad14-058f828fd097" +``` + +**Remove all datasets from a specific collection (careful!):** +```bash +owi local rm "all/collectionName=test_collection" +``` + +--- + +## `local analyze` + +Analyze the content of local json files created with a queyr command (e.g. `query less`) without treating them as full datasets. + +### Usage + +```bash +owi local analyze [FILES]... [OPTIONS] +``` + +### Options + +- `--topk N`: Number of top frequent items to show. +- `--sample N`: Number of lines/rows to sample. + +### Examples + +**Analyze a JSONL file:** +```bash + owi --format json query less -L all/id=0350fecc-e58b-11f0-a8c9-8ebf6bb2cab9 >~/tmp/test.json +owi local analyze ~/tmp/test.json +``` diff --git a/docs/source/details/query.md b/docs/source/details/query.md new file mode 100644 index 0000000..a11b372 --- /dev/null +++ b/docs/source/details/query.md @@ -0,0 +1,193 @@ +# Query Commands + +Query commands allow you to perform SQL-based analysis and filtering on datasets (typically Parquet files) without downloading the entire dataset. They use DuckDB for efficient querying. + +## `query less` + +An interactive browser for dataset content, similar to the unix `less` command but for structured data. + +### Usage + +```bash +owi query less [OPTIONS] +``` + +### Options + +- `-R SPECIFIER` / `--remote`: Remote dataset specifier (e.g., `lexis:latest`). +- `-L SPECIFIER` / `--local`: Local dataset specifier (e.g., `all:latest`). +- `--select SQL`: SQL selection string (default: `url,title,main_content`). +- `--where SQL`: SQL where clause for filtering. +- `--limit N`: Limit number of rows. +- `--page-size N`: Rows per page in table view (default: 10). +- `--pq-batch N`: Number of parquet files per query batch (default: 10). +- `--batch-size N`: Rows per processing batch (default: 100). +- `--files GLOB`: Glob pattern to select files (default: `**/*.parquet`). +- `--async/--sync`: Execution mode (see below). + +### Execution Modes + +- **Sync mode (default)**: Results stream as they arrive. Good for interactive use and seeing results quickly. +- **Async mode (`--async`)**: Uses concurrent async execution which can be faster for large datasets, but results stream via a background thread (overhead) + +### Examples + +**Interactive table view of remote data:** +```bash +owi query less -R lexis:latest --select "url,title" +``` + +**Stream results as JSONL with progress:** +```bash +owi --format jsonl+ query less -R lexis:latest --limit 1000 +``` + +**Query local dataset:** +```bash +owi query less -L all:latest/collectionName=main --limit 100 +``` + +**High-performance async mode:** +```bash +owi query less -R lexis:latest --async --pq-batch 10 +``` + +--- + +## `query stats` + +Calculate statistics for a dataset, such as record counts, date ranges, and distributions. + +### Usage + +```bash +owi query stats [OPTIONS] +``` + +### Options + +- `-R SPECIFIER` / `--remote`: Remote dataset specifier. +- `-L SPECIFIER` / `--local`: Local dataset specifier. +- `--topk N`: Number of top frequent items to show for distributions. +- `--output FILE`: Save statistics to a JSON/text file. + +### Examples + +**Get statistics for a remote dataset:** +```bash +owi query stats -R lexis:latest +``` + +--- + +## `query sites` + +Filter datasets based on a list of URLs or domains provided in a file. + +### Usage + +```bash +owi query sites [OPTIONS] URL_FILE +``` + +### Options + +- `-R SPECIFIER` / `--remote`: Remote dataset specifier. +- `-L SPECIFIER` / `--local`: Local dataset specifier. +- `--select SQL`: Columns to select. +- `--output FILE`: Output file for results. + +### Examples + +**Extract records for domains listed in `domains.txt`:** +```bash +owi query sites -R lexis:latest/collectionName=curlie_full domains.txt --output filtered_data.jsonl +``` + +--- + +## `query aggregate` + +Run arbitrary SQL aggregation queries (GROUP BY) on datasets. + +### Usage + +```bash +owi query aggregate [OPTIONS] +``` + +### Options + +- `-R SPECIFIER` / `--remote`: Remote dataset specifier. +- `-L SPECIFIER` / `--local`: Local dataset specifier. +- `--select SQL`: Selection including aggregation functions (e.g., `count(*) as c`). +- `--groupby SQL`: Columns to group by. +- `--aggregate SQL`: Final logic on the aggregated result (e.g. `ORDER BY c DESC`). + +### Examples + +**Count URLs by top-level domain:** +```bash +owi query aggregate -R lexis:latest/collectionName=curlie_full --select "url_suffix, count(*) as c" --groupby "url_suffix" +``` + +--- + +### Subcommands + +#### `owi query aggregate` +Run SQL aggregation queries across datasets. Automatically merges partial results for count/sum. + +```bash +owi query aggregate -R lexis:latest [OPTIONS] +``` +- `-s, --select TEXT`: SQL SELECT (default: "url_suffix,count(url) as c") +- `-g, --groupby TEXT`: SQL GROUP BY (default: "url_suffix") +- `-a, --aggregate TEXT`: Final aggregation query +- `--async`: Use asynchronous execution (streaming) + +#### `owi query slice` +Create a new dataset slice based on a query. See [Query Slice Details](query_slice.md). + +```bash +owi query slice -R lexis:latest/collectionName=source [OPTIONS] +``` +- `-c, --collection TEXT`: Name for the new slice collection (default: "userslice") +- `-s, --select TEXT`: Columns to include (default: "*") +- `-w, --where TEXT`: Filter condition +- `-l, --limit TEXT`: Row limit +- `-t, --type TEXT`: Resource type (default: "owip") +- `-a, --access TEXT`: Access level (default: "public") +- `--ciff/--no-ciff`: Also slice associated CIFF files (default: True) +- `--overwrite`: Overwrite existing slice + +#### `owi query stats` +Calculate statistics over datasets (counts, distributions). + +```bash +owi query stats -R lexis:latest [OPTIONS] +``` +- `-k, --topk INTEGER`: Top N items for distributions (default: 20) +- `-w, --write-back`: Save stats to dataset .stats folder +- `--force`: Force recalculation +- `--async`: Use asynchronous execution + +#### `owi query sites` +Execute site-specific queries using a URL file filter. + +```bash +owi query sites urls.txt -R lexis:latest [OPTIONS] +``` +- `-s, --select TEXT`: Columns to select +- `-w, --where TEXT`: Additional filter +- `--async`: Use asynchronous execution + +#### `owi query warc` +Extract WARC file locations from datasets. + +```bash +owi query warc -R lexis:latest [OPTIONS] +``` +- `-w, --where TEXT`: Filter condition +- `--warc-config PATH`: WARC location config file +- `-o, --output PATH`: Output directory for logs diff --git a/docs/source/details/query_slice.md b/docs/source/details/query_slice.md new file mode 100644 index 0000000..7b6196c --- /dev/null +++ b/docs/source/details/query_slice.md @@ -0,0 +1,98 @@ +# Query Slice Command + +The `query slice` command is a powerful tool for creating new datasets (subsets) from existing local or remote datasets based on SQL queries. It leverages DuckDB to efficiently process parquet files and generate new datasets adhering to the OWI metadata schema. + +## Usage + +```bash +owi query slice [OPTIONS] +``` + +### Key Arguments + +- **Target Source**: + - `-L, --local SPECIFIER`: Slice from a local dataset. + - `-R, --remote SPECIFIER`: Slice from a remote dataset (requires files to be accessible/pulled or creates empty structure for metadata slicing). *Note: Slicing usually operates on downloaded files.* + +- **Filtering & Selection**: + - `--where TEXT`: SQL WHERE clause to filter records (e.g., `url_suffix='de'`). + - `--select TEXT`: SQL SELECT clause to choose columns (default: `*`). + +- **Output Dataset**: + - `--collection NAME`: Name of the collection for the new dataset. + - `--access LEVEL`: Access level for the new dataset (default: `public`). Options: `public`, `project`, `user`. + - `--id UUID`: Target an **existing dataset** by ID. If specified, the slice result is appended/merged into this dataset. + +- **Execution Control**: + - `--yes`: Skip confirmation prompts. + - `--format FORMAT`: Output format (`table`, `json`, `jsonl`). + +## Workflows + +### 1. Creating a New Slice + +To create a new dataset containing only German URLs from a source dataset: + +```bash +owi query slice -L "all:latest" --where "url_suffix='de'" --collection "german_subset" --yes +``` + +This will: +1. Discover files in the source dataset. +2. Execute the DuckDB query `SELECT * FROM source WHERE url_suffix='de'`. +3. Write the result to a new dataset under `~/.owi/public/german_subset/`. +4. Infer metadata (file counts, size) and save it. + +### 2. Updating (Appending to) a Slice + +To add French URLs to the **same** dataset created above, use the `--id` flag: + +```bash +# First, find the ID of the dataset created in step 1 +owi local ls all/collectionName=german_subset + +# Then slice and append +owi query slice -L "all:latest" --where "url_suffix='fr'" --collection "german_subset" --id --yes +``` + +### 3. Automated Slicing with JSON Output + +For scripts and automation pipelines, use `--format json`. This outputs a structured JSON array containing details of the created/updated dataset, while suppressing progress bars. + +```bash +owi query slice -L "all:latest" ... --format json +``` + +**Output Example:** + +```json +[ + { + "dataset": { + "id": "15c682b7-7615-470e-8669-02fc667e3a62", + "path": "/Users/user/.owi/public/german_subset/15c682b7-...", + "collectionName": "german_subset", + "access": "public", + "title": "OWI-Open Web Index-german_subset...", + "files": 150, + "size": 104857600 + } + } +] +``` + +## Technical Details + +### Access Levels +The `access` parameter determines where the dataset is stored locally: +- `public`: `~/.owi/public//` +- `project`: `~/.owi/project//` +- `user`: `~/.owi/user//` + +*Fix Note*: Previous issues with `access` being ignored (defaulting to `None`) have been resolved. + +### Data Provenance +The sliced dataset includes provenance information in its metadata, recording the source dataset ID, the query used (`select`, `where`), and the timestamp of the slice operation. + +### CIFF Files +If the `--slice-ciff` flag is enabled (default), the command also attempts to slice associated CIFF (Common Index File Format) files if they exist and match the query filters. diff --git a/docs/source/details/remote.md b/docs/source/details/remote.md new file mode 100644 index 0000000..c91d44d --- /dev/null +++ b/docs/source/details/remote.md @@ -0,0 +1,145 @@ +# Remote Commands + +Remote commands (`owi remote`) allow you to interact with datasets stored in remote data centers (e.g., LEXIS, LRZ, IT4I). + +## `remote ls` + +List datasets available in remote repositories. + +### Usage + +```bash +owi remote ls [SPECIFIER] [OPTIONS] +``` + +### Options + +- `--files GLOB`: List files within the remote dataset(s) matching the glob pattern. Note that this requires fetching file lists from the remote server. +- `--file-details`: When used with `--files`, shows file sizes (slower as it fetches metadata for each file). +- `--display FORMAT`: Output format (`short`, `mid`, `wide`, `full`). + +### Examples + +**List the latest datasets from the 'lexis' datacenter:** +```bash +owi remote ls lexis:latest +``` + +**List all parquet files in a remote dataset:** +```bash +owi remote ls lexis:latest/collectionName=main --files "**/*.parquet" +``` + +**List files with sizes (slower):** +```bash +owi remote ls lexis:latest/collectionName=main --files "**/*.parquet" --file-details +``` + +--- + +## `remote pull` + +Download datasets or specific files from remote repositories to your local environment. + +### Usage + +```bash +owi remote pull [SPECIFIER] [OPTIONS] +``` + +### Options + +- `--files GLOB`: Only download files matching the glob pattern. +- `--overwrite`: Overwrite existing local files. +- `--threads N`: Number of concurrent download threads. +- `--language LANG`: Filter files by language code (if applicable to file structure). +- `--day DATE`: Filter by specific day (snapshot). + +### Examples + +**Pull an entire dataset:** +```bash +owi remote pull lexis:latest +``` + +**Pull only parquet files:** +```bash +owi remote pull lexis:latest --files "**/*.parquet" +``` + +**Parallel download:** +```bash +owi remote pull lexis:latest --threads 8 +``` + +--- + +## `remote push` + +Upload local datasets to a remote repository. + +### Usage + +```bash +owi remote push [SPECIFIER] [OPTIONS] +``` + +### Options + +- `--datacenter DC`: Target datacenter (if not specified in specifier). +- `--overwrite`: Overwrite remote files if they exist. +- `--yes` (`-y`): Skip confirmation. + +### Examples + +**Push a local dataset to IT4I:** +```bash +owi remote push "all/id=my-local-id" --datacenter it4i +``` + +--- + +## `remote diff` + +Compare datasets between local and remote repositories to see what is missing or different. + +### Usage + +```bash +owi remote diff [SPECIFIER] [OPTIONS] +``` + +### Options + +- `--files GLOB`: Perform file-level comparison using the glob pattern. + +### Examples + +**Compare local datasets with remote 'project' access datasets:** +```bash +owi remote diff "all/access=project" +``` + +--- + +## `remote doctor` + +Check the connection status and authentication to configured remote data centers. + +### Usage + +```bash +owi remote doctor +``` + +--- + +## `remote logout` + +Revoke the current access token, effectively logging out from remote services. + +### Usage + +```bash +owi remote logout +``` diff --git a/docs/source/details/warc.md b/docs/source/details/warc.md index 997fc68..68e425b 100644 --- a/docs/source/details/warc.md +++ b/docs/source/details/warc.md @@ -1,7 +1,9 @@ # WARC Cache and Download Module ```{eval-rst} -.. automodule:: owilix.cmd.subcmds.query_warc +.. automodule:: owilix.core.tasks.warc + :members: + :undoc-members: :show-inheritance: ``` @@ -14,35 +16,35 @@ Screenshot of the warc module usage: ## WARC Download function (used as OWILIX query subcommand) ```{eval-rst} -.. autofunction:: owilix.cmd.subcmds.query_warc.warc +.. autofunction:: owilix.core.tasks.warc.query_warc.warc ``` ## Log Analysis Function ```{eval-rst} -.. autofunction:: owilix.cmd.subcmds.query_warc.analyze_warc_log +.. autofunction:: owilix.core.tasks.warc.parquet_logger.analyze_warc_log ``` # Module Components and Classes ```{eval-rst} -.. autoclass:: owilix.cmd.subcmds.query_warc.query_warc.ZMQStreamingWARCProcessor +.. autoclass:: owilix.core.tasks.warc.query_warc.ZMQStreamingWARCProcessor :members: :special-members: __init__ ``` -### Per-Thread Destination System +## Per-Thread Destination System The core innovation is the per-thread destination architecture: ```{eval-rst} -.. autoclass:: owilix.cmd.subcmds.query_warc.query_warc.ParallelWARCDestinationManager +.. autoclass:: owilix.core.tasks.warc.query_warc.ParallelWARCDestinationManager :members: ``` ```{eval-rst} -.. autoclass:: owilix.cmd.subcmds.query_warc.query_warc.WARCDestination +.. autoclass:: owilix.core.tasks.warc.query_warc.WARCDestination :members: ``` @@ -51,12 +53,12 @@ The core innovation is the per-thread destination architecture: ### Metrics Collection ```{eval-rst} -.. autoclass:: owilix.cmd.subcmds.query_warc.query_warc.DatacenterStats +.. autoclass:: owilix.core.tasks.warc.query_warc.DatacenterStats :members: ``` ```{eval-rst} -.. autoclass:: owilix.cmd.subcmds.query_warc.query_warc.WARCDestinationStats +.. autoclass:: owilix.core.tasks.warc.query_warc.WARCDestinationStats :members: ``` @@ -66,12 +68,12 @@ The core innovation is the per-thread destination architecture: ### Parquet-Based Job Logging ```{eval-rst} -.. autoclass:: owilix.cmd.subcmds.query_warc.parquet_logger.ParquetJobLogger +.. autoclass:: owilix.core.tasks.warc.parquet_logger.ParquetJobLogger :members: ``` ```{eval-rst} -.. autoclass:: owilix.cmd.subcmds.query_warc.parquet_logger.JobLogEntry +.. autoclass:: owilix.core.tasks.warc.parquet_logger.JobLogEntry :members: ``` @@ -81,7 +83,7 @@ The core innovation is the per-thread destination architecture: ### File Processor ```{eval-rst} -.. autoclass:: owilix.cmd.subcmds.query_warc.query_warc.HighPerformanceFileProcessor +.. autoclass:: owilix.core.tasks.warc.query_warc.HighPerformanceFileProcessor :members: ``` diff --git a/docs/source/index.rst b/docs/source/index.rst index 2f4bb46..30b0939 100644 --- a/docs/source/index.rst +++ b/docs/source/index.rst @@ -22,7 +22,27 @@ Welcome to Owilix's documentation! troubleshooting.md dev.md contributing/index + details/local.md + details/remote.md + details/query.md + details/config.md details/warc.md + details/query_slice.md + details/graphs.md + +.. toctree:: + :maxdepth: 1 + :caption: Developer Reference: + + config.md + core.md + db_integration.md + fsspec_integration.md + migration.md + repository.md + server.md + streams.md + testing_and_benchmarks.md .. automodule:: owilix :members: diff --git a/owilix/cli/local.py b/owilix/cli/local.py index c53ea7f..0652167 100644 --- a/owilix/cli/local.py +++ b/owilix/cli/local.py @@ -23,6 +23,17 @@ app = typer.Typer( ) +def _format_size(size_bytes: int) -> str: + """Format bytes as human-readable string.""" + if size_bytes < 1024: + return f"{size_bytes} B" + elif size_bytes < 1024 ** 2: + return f"{size_bytes / 1024:.1f} KiB" + elif size_bytes < 1024 ** 3: + return f"{size_bytes / 1024 ** 2:.1f} MiB" + else: + return f"{size_bytes / 1024 ** 3:.2f} GiB" + @app.command() def ls( ctx: typer.Context, @@ -33,6 +44,7 @@ def ls( no_summary: bool = typer.Option(False, "--no-summary", help="Skip summary table"), fields: Optional[str] = typer.Option(None, "--fields", help="Customize fields (+field, -field)"), files_glob: Optional[str] = typer.Option(None, "--files", "-f", help="List files matching glob pattern (e.g., '**/*.parquet')"), + file_details: bool = typer.Option(False, "--file-details", help="Show file details (size) - slower"), ): """ List local datasets matching SPECIFIER. @@ -63,22 +75,41 @@ def ls( # If --files is specified, list files instead of datasets if files_glob: from rich.table import Table + from rich import box for ds in datasets_list: cli_ctx.console.print(f"\n[bold]Dataset:[/bold] {ds.metadata.get('title', 'Unknown')} ({ds.metadata.id})") cli_ctx.console.print(f"[dim]Path: {ds.path}[/dim]") - files = ds.files(files_glob, relative=True) - if files: - table = Table(title=f"Files matching '{files_glob}'", show_header=True) - table.add_column("File", style="cyan") - for f in sorted(files)[:50]: # Limit to 50 files - table.add_row(f) - if len(files) > 50: - table.add_row(f"... and {len(files) - 50} more files") - cli_ctx.console.print(table) - cli_ctx.console.print(f"[green]Total: {len(files)} files[/green]") + if file_details: + # Use files_details to get size and other metadata (slower) + details = ds.files_details(files_glob, count_rows=False) + if details: + table = Table(title=f"Files matching '{files_glob}'", show_header=True, box=box.SIMPLE) + table.add_column("File", style="cyan") + table.add_column("Size", style="dim", justify="right") + + total_size = 0 + for f in sorted(details, key=lambda x: x.get("relpath", x.get("path", ""))): + size = f.get("info", {}).get("size", 0) + total_size += size + table.add_row(f.get("relpath", f.get("path", "?")), _format_size(size)) + + cli_ctx.console.print(table) + cli_ctx.console.print(f"[green]Total: {len(details)} files, {_format_size(total_size)}[/green]") + else: + cli_ctx.console.print(f"[yellow]No files matching '{files_glob}'[/yellow]") else: - cli_ctx.console.print(f"[yellow]No files matching '{files_glob}'[/yellow]") + # Fast mode: just list file names + files = ds.files(files_glob, relative=True) + if files: + table = Table(title=f"Files matching '{files_glob}'", show_header=True, box=box.SIMPLE) + table.add_column("File", style="cyan") + for f in sorted(files): + table.add_row(f) + cli_ctx.console.print(table) + cli_ctx.console.print(f"[green]Total: {len(files)} files[/green]") + else: + cli_ctx.console.print(f"[yellow]No files matching '{files_glob}'[/yellow]") return # Sort if requested @@ -135,7 +166,7 @@ def rm( # Parse specifier and list datasets spec = cli_ctx.owi.parse_specifier(specifier) datasets = cli_ctx.owi.local.list( - access=spec.get("query", {}).get("access"), + access=spec.get("query", {}).get("access", "public"), query={k: v for k, v in spec.get("query", {}).items() if k != "access"}, ) @@ -166,8 +197,8 @@ def rm( progress.update(task, current_item=title[:30]) try: - cli_ctx.owi.local.remove(ds) - cli_ctx.console.print(f"[green]✓[/green] Removed {title}") + cli_ctx.owi.local.delete(ds) + cli_ctx.console.print(f"[green]Removed {ds.metadata.get('title', 'Unknown')}[/green]") except Exception as e: cli_ctx.console.print(f"[red]✗[/red] Failed to remove {title}: {e}") @@ -199,13 +230,16 @@ def analyze( local_cmd = LocalCommands( cli_ctx.owi, - console=cli_ctx.console, + CONSOLE=cli_ctx.console, + TARGET=cli_ctx.target, + AUTOYES=cli_ctx.auto_yes, + owilix_config=cli_ctx.owi.config, ) - result = local_cmd.analyze_jsonl(path, topk=topk) + res = local_cmd.commands["analyze_jsonl"](path, topk=topk) - if result and 'markdown' in result: - cli_ctx.console.print(result['markdown']) + if res.success and res.object and 'markdown' in res.object: + cli_ctx.console.print(res.object['markdown']) @app.command() diff --git a/owilix/cli/query.py b/owilix/cli/query.py index 77e66f1..1514d68 100644 --- a/owilix/cli/query.py +++ b/owilix/cli/query.py @@ -20,6 +20,7 @@ from owilix.core.tasks.query import ( slice_query, stream_query, query_stats ) from owilix.core.tasks.warc.query_warc import warc as query_warc +from owilix.cli._common.output import OutputWriter # Create sub-app for query commands app = typer.Typer( @@ -44,7 +45,7 @@ def less( verbose: bool = typer.Option(False, "--verbose", help="Enable verbose output"), job_name: Optional[str] = typer.Option(None, "--job", help="Job name for transaction logging"), resume: bool = typer.Option(False, "--resume", help="Resume from previous execution"), - async_mode: bool = typer.Option(True, "--async/--sync", help="Use async executor"), + async_mode: bool = typer.Option(False, "--async/--sync", help="Use async executor (sync streams results, async collects all first)"), ): """ Interactive dataset browser with SQL query capabilities. @@ -104,6 +105,7 @@ def sites( verbose: bool = typer.Option(False, "--verbose", help="Enable verbose output"), job_name: Optional[str] = typer.Option(None, "--job", help="Job name for transaction logging"), resume: bool = typer.Option(False, "--resume", help="Resume from previous execution"), + async_mode: bool = typer.Option(False, "--async/--sync", help="Use async executor"), ): """ Execute site-specific queries using URL filtering. @@ -139,6 +141,7 @@ def sites( verbose=verbose, job_name=job_name, resume=resume, + async_mode=async_mode, command_name="query sites" ) @@ -193,6 +196,7 @@ def aggregate( batch_size: int = typer.Option(100, "--batch-size", help="Rows per batch"), pq_batch_size: int = typer.Option(1, "--pq-batch", help="Parquet files per batch"), verbose: bool = typer.Option(False, "--verbose", help="Enable verbose output"), + async_mode: bool = typer.Option(False, "--async/--sync", help="Use async executor"), ): """ Run SQL aggregation queries across datasets. @@ -215,6 +219,7 @@ def aggregate( output_file=cli_ctx.output_file, batch_size=batch_size, pq_batch_size=pq_batch_size, + async_mode=async_mode, console=cli_ctx.console, command_name="query aggregate" ) @@ -232,6 +237,7 @@ def stats( topk: int = typer.Option(20, "--topk", "-k", help="Number of top items per category"), write_back: bool = typer.Option(False, "--write-back", "-w", help="Save stats to dataset .stats folder"), force_write: bool = typer.Option(False, "--force", help="Force overwrite existing stats"), + async_mode: bool = typer.Option(False, "--async/--sync", help="Use async executor"), ): """ Calculate statistics over datasets. @@ -258,6 +264,7 @@ def stats( output_file=cli_ctx.output_file, write_back=write_back, force_write=force_write, + async_mode=async_mode, console=cli_ctx.console, command_name="query stats" ) @@ -292,11 +299,16 @@ def slice( resume: bool = typer.Option(False, "--resume", help="Resume previous execution"), verbose: bool = typer.Option(False, "--verbose", help="Enable verbose output"), yes: bool = typer.Option(False, "--yes", "-y", help="Skip confirmation"), + # Output format options using global enums from main + format: Optional[str] = typer.Option(None, "--format", help="Output format (table, json, jsonl, json+, jsonl+)"), + output_file: Optional[str] = typer.Option(None, "--output", "-o", help="File to write output to"), ): """ Create a new dataset by slicing data from existing datasets. """ cli_ctx: CLIContext = ctx.obj + # Resolve format: command-line arg > global arg > default + format = format or cli_ctx.output_format or "table" # Perform slice directly (as it prints progress etc extensively) result = slice_query( @@ -323,7 +335,7 @@ def slice( job_name=job_name, resume=resume, verbose=verbose, - console=cli_ctx.console, + console=cli_ctx.console if format == "table" and not output_file else None, autoyes=yes, ) @@ -331,7 +343,20 @@ def slice( raise typer.Exit(code=1) if result and result.success: - cli_ctx.console.print(f"[green]✓ Slice completed successfully[/green]") + if format in ("json", "jsonl", "json+", "jsonl+") or output_file: + # Structured output for automation + output_data = result.object + if not output_data: + output_data = {"status": "success", "message": "Slice completed successfully"} + + # Update context for OutputWriter + cli_ctx.output_format = format + cli_ctx.output_file = output_file + + with OutputWriter(cli_ctx) as writer: + writer.write_records([output_data]) + else: + cli_ctx.console.print(f"[green]✓ Slice completed successfully[/green]") @app.command() @@ -353,11 +378,15 @@ def stream( prefetch: int = typer.Option(1, "--prefetch", help="Batches to prefetch"), queue_size: int = typer.Option(5, "--queue", help="Buffer queue size"), verbose: bool = typer.Option(False, "--verbose", help="Enable verbose output"), + format: Optional[str] = typer.Option(None, "--format", help="Output format (table, json, jsonl, json+, jsonl+)"), + output_file: Optional[str] = typer.Option(None, "--output", "-o", help="File to write output to"), ): """ Stream query results to stdout or custom consumer. """ cli_ctx: CLIContext = ctx.obj + # Resolve format: command-line arg > global arg > default + format = format or cli_ctx.output_format or "table" result = stream_query( owi=cli_ctx.owi, diff --git a/owilix/cli/remote.py b/owilix/cli/remote.py index 793bec2..dddff00 100644 --- a/owilix/cli/remote.py +++ b/owilix/cli/remote.py @@ -22,6 +22,18 @@ app = typer.Typer( ) +def _format_size(size_bytes: int) -> str: + """Format bytes as human-readable string.""" + if size_bytes < 1024: + return f"{size_bytes} B" + elif size_bytes < 1024 ** 2: + return f"{size_bytes / 1024:.1f} KiB" + elif size_bytes < 1024 ** 3: + return f"{size_bytes / 1024 ** 2:.1f} MiB" + else: + return f"{size_bytes / 1024 ** 3:.2f} GiB" + + @app.command() def ls( ctx: typer.Context, @@ -32,6 +44,7 @@ def ls( no_summary: bool = typer.Option(False, "--no-summary", help="Skip summary table"), fields: Optional[str] = typer.Option(None, "--fields", help="Customize fields (+field, -field)"), files_glob: Optional[str] = typer.Option(None, "--files", "-f", help="List files matching glob pattern (e.g., '**/*.parquet')"), + file_details: bool = typer.Option(False, "--file-details", help="Show file details (size) - slower"), ): """ List remote datasets matching SPECIFIER. @@ -65,23 +78,42 @@ def ls( # If --files is specified, list files instead of datasets if files_glob: from rich.table import Table + from rich import box for ds in datasets_list: cli_ctx.console.print(f"\n[bold]Dataset:[/bold] {ds.metadata.get('title', 'Unknown')} ({ds.metadata.id})") cli_ctx.console.print(f"[dim]DataCenter: {getattr(ds, 'dataCenter', 'unknown')}[/dim]") try: - files = ds.files(files_glob, relative=True) - if files: - table = Table(title=f"Files matching '{files_glob}'", show_header=True) - table.add_column("File", style="cyan") - for f in sorted(files)[:50]: # Limit to 50 files - table.add_row(f) - if len(files) > 50: - table.add_row(f"... and {len(files) - 50} more files") - cli_ctx.console.print(table) - cli_ctx.console.print(f"[green]Total: {len(files)} files[/green]") + if file_details: + # Use files_details to get size and other metadata (slower) + details = ds.files_details(files_glob, count_rows=False) + if details: + table = Table(title=f"Files matching '{files_glob}'", show_header=True, box=box.SIMPLE) + table.add_column("File", style="cyan") + table.add_column("Size", style="dim", justify="right") + + total_size = 0 + for f in sorted(details, key=lambda x: x.get("relpath", x.get("path", ""))): + size = f.get("info", {}).get("size", 0) + total_size += size + table.add_row(f.get("relpath", f.get("path", "?")), _format_size(size)) + + cli_ctx.console.print(table) + cli_ctx.console.print(f"[green]Total: {len(details)} files, {_format_size(total_size)}[/green]") + else: + cli_ctx.console.print(f"[yellow]No files matching '{files_glob}'[/yellow]") else: - cli_ctx.console.print(f"[yellow]No files matching '{files_glob}'[/yellow]") + # Fast mode: just list file names + files = ds.files(files_glob, relative=True) + if files: + table = Table(title=f"Files matching '{files_glob}'", show_header=True, box=box.SIMPLE) + table.add_column("File", style="cyan") + for f in sorted(files): + table.add_row(f) + cli_ctx.console.print(table) + cli_ctx.console.print(f"[green]Total: {len(files)} files[/green]") + else: + cli_ctx.console.print(f"[yellow]No files matching '{files_glob}'[/yellow]") except Exception as e: cli_ctx.console.print(f"[red]Error listing files: {e}[/red]") return @@ -126,41 +158,22 @@ def ls( @app.command() def doctor( ctx: typer.Context, + tokens: bool = typer.Option(False, "--tokens", "-t", help="Show session token information"), ): """ Check connection status of configured remotes. + + Use --tokens to display session token information (truncated for security). """ cli_ctx: CLIContext = ctx.obj + from owilix.core.tasks.remote import remote_doctor - # Get repos and check connectivity - repos = cli_ctx.owi.remote_data.get_repos() - - with OutputWriter(cli_ctx) as writer: - if cli_ctx.output_format == "table": - records = [] - for name, repo in repos.items(): - # Check if repo has is_connected, otherwise just show as configured - if hasattr(repo, 'is_connected'): - connected = repo.is_connected() - status = "✅ connected" if connected else "❌ disconnected" - else: - status = "⚪ configured" - records.append({ - "name": name, - "type": type(repo).__name__, - "status": status, - }) - writer.write_records(records, title="Remote Repositories") - else: - records = [] - for name, repo in repos.items(): - connected = repo.is_connected() if hasattr(repo, 'is_connected') else None - records.append({ - "name": name, - "type": type(repo).__name__, - "connected": connected - }) - writer.write_records(records) + remote_doctor( + manager=cli_ctx.owi, + as_json=cli_ctx.output_format in ("json", "jsonl"), + show_tokens=tokens, + console=cli_ctx.console, + ) @app.command() diff --git a/owilix/core/db/duckdb_executor.py b/owilix/core/db/duckdb_executor.py index 7a0a5c8..598e008 100644 --- a/owilix/core/db/duckdb_executor.py +++ b/owilix/core/db/duckdb_executor.py @@ -109,7 +109,7 @@ class OWIDuckDBSelectExecutor: :param output_format: output format to yield results in. values are "dict" or "list" or "arrow" :param fn_group: A function that groups the (file_path, prefix) list into ParquetBatch objects """ - self.logger = logging.getLogger("OWIDuckDBSelectExecutor") + self.logger = logging.getLogger("owilix.db") self.sql = owilix_sql self.pq_batch_size = pq_batch_size self.batch_size = batch_size @@ -254,9 +254,10 @@ class OWIDuckDBSelectExecutor: """ conn_tuple = None try: - + self.logger.debug(f"run_query_on_batch: acquiring connection for {len(pq_batch.files)} files") conn_tuple = self.pool.acquire_connection() conn, tmp_dir = conn_tuple + self.logger.debug(f"run_query_on_batch: connection acquired") # Register the filesystem with the connection if needed # fs.protocol can be a string or tuple (e.g., ('file', 'local')) @@ -274,17 +275,19 @@ class OWIDuckDBSelectExecutor: # Prepare the query file_list = [f[0] for f in pq_batch.files] + self.logger.debug(f"run_query_on_batch: building query for {len(file_list)} files") _query = self.sql.files(file_list) if file_list else self.sql # If we have query_args in the batch, format the SQL if pq_batch.query_args: _query = _query.format(**pq_batch.query_args) - self.logger.debug(f"Executing query: {_query.sql} on files={file_list}") + self.logger.debug(f"run_query_on_batch: executing query on {len(file_list)} files...") # Execute with retry try: cursor = self._retry_query(conn.execute, _query.sql, fs=fs) + self.logger.debug(f"run_query_on_batch: query executed, fetching results...") except Exception as e: # If a fetch fails, we can decide to re-run the entire query from scratch self.logger.exception(f"Entire Query failed. {e}") @@ -390,7 +393,7 @@ class OWIDuckDBSelectExecutor: as they come in. """ tasks = self.generate_ordered_tasks() - self.logger.info(f"Scheduling {len(tasks)} tasks...") + self.logger.info(f"query_aggregator: scheduling {len(tasks)} tasks with prefetch={self.prefetch}") with ThreadPoolExecutor(max_workers=self.prefetch) as executor: future_to_task = { @@ -413,6 +416,52 @@ class OWIDuckDBSelectExecutor: error=e ) + async def query_aggregator_async(self): + """ + Async version of query_aggregator that runs queries in a thread pool. + + This allows non-blocking consumption of query results from async code. + """ + import asyncio + from concurrent.futures import ThreadPoolExecutor + + tasks = self.generate_ordered_tasks() + self.logger.info(f"query_aggregator_async: scheduling {len(tasks)} tasks with prefetch={self.prefetch}") + + loop = asyncio.get_event_loop() + + with ThreadPoolExecutor(max_workers=self.prefetch) as executor: + # Submit all tasks + futures = [ + loop.run_in_executor(executor, self._run_query_on_batch_collect, fs, batch) + for fs, batch in tasks + ] + + # Yield results as they complete + for future in asyncio.as_completed(futures): + try: + results = await future + for result in results: + yield result + except Exception as e: + self.logger.exception(f"Async task failed: {e}") + yield ParquetBatchResult( + parquet_batch=ParquetBatch(files=[], query_args={}), + success=False, + error=e + ) + + def _run_query_on_batch_collect(self, fs, batch): + """Collect all results from run_query_on_batch generator into a list.""" + return list(self.run_query_on_batch(fs, batch)) + + async def close(self): + """Close the executor and cleanup resources.""" + try: + self.pool.close_all() + except: + pass + class OWIDuckDBCopyExecutor(OWIDuckDBSelectExecutor): """ Executor that implements a COPY operation with chunking. diff --git a/owilix/core/models/dataset.py b/owilix/core/models/dataset.py index 3946517..f4688f6 100644 --- a/owilix/core/models/dataset.py +++ b/owilix/core/models/dataset.py @@ -4,6 +4,7 @@ import copy import os import platform import re +import logging from collections import defaultdict from datetime import datetime, timedelta, date from typing import Literal @@ -82,7 +83,9 @@ def infer_metadata_from_files(files, cb = None): _p = f["relpath"] if isinstance(f, dict) else f if cb is not None: cb(_p, ix, len(f)) - if len(_p.split("/"))<2: continue + # Ensure we don't skip files in root (e.g. part-0.parquet) + # if len(_p.split("/"))<2: continue + if any([n.startswith(".") for n in _p.split("/")]): continue _files.append(_p) if isinstance(f, dict): @@ -104,6 +107,22 @@ def infer_metadata_from_files(files, cb = None): if any([".warc" in f for f in _files]): _returns["subResourceType"].append("warc") groups = group_and_count_files(_files, "", 20) _dates = set([parser.parse(p) for p in _extract_unique_dates(set(groups.keys()))]) + + # Merge with dates from file content metadata + for f in files: + if isinstance(f, dict): + info = f.get("info", {}) + for key in ["startDate", "endDate"]: + if key in info and info[key]: + try: + d = info[key] + if isinstance(d, (str, bytes)): + _dates.add(parser.parse(str(d))) + elif hasattr(d, 'year'): # datetime/date object + _dates.add(d) + except: + pass + if len(_dates)>0: _returns["startDate"] = min(_dates) _returns["endDate"] = max(_dates) @@ -1133,7 +1152,10 @@ class AdditionalMetadata(MetadataField): _cast = int if v["key"] in ["fileCount","objectCount", "totalSize"] else lambda x: x self.additionalMetadata[v["key"]]=_cast(v["value"]) else: - self.additionalMetadata = self._cast(value) if isinstance(value, dict) else {} + if isinstance(value, dict): + self.additionalMetadata.update(self._cast(value)) + else: + self.additionalMetadata = {} except: return def _cast(self, _d:dict): @@ -1151,7 +1173,7 @@ class AdditionalMetadata(MetadataField): def to_json_dict(self, version: str): """returns a dictionary that is json serializabel for the version provided. """ - return {"additionalMetadata": self.to_value(version)} + return self.to_value(version) def short_repr(self): values = '-'.join(str(v) for v in self.additionalMetadata.values()) @@ -1553,7 +1575,10 @@ class DatasetMetadata: if isinstance(entry, MetadataField): value = entry.to_value(self.default_version) if self.default_version == 'datacite.V1': - datacite_result.update(value) + if key == "additionalMetadata": + result["additionalMetadata"] = entry.to_value("V1") + else: + datacite_result.update(value) else: result.update(value) else: @@ -1984,6 +2009,7 @@ class Dataset: """ Create a new dataset by merging the provided ones and applying files, select and where filters (for updating the provenance) """ + logging.getLogger(__name__).debug(f"merge_into_new called with access={access}, collectionName={collectionName}") _new_md = DatasetMetadata.new_metadata(resourceType, creator ,**kwargs) _new_md["provenance"] = [f"{create_provenance_url(d, files, select=select, where=where)}" for d in datasets] diff --git a/owilix/core/tasks/query.py b/owilix/core/tasks/query.py index f697b69..c8e94ec 100644 --- a/owilix/core/tasks/query.py +++ b/owilix/core/tasks/query.py @@ -156,6 +156,7 @@ def query_sites( verbose: bool = False, job_name: Optional[str] = None, resume: bool = False, + async_mode: bool = False, console: Optional[Console] = None ) -> CommandResult: """Execute site-specific queries using URL filtering.""" @@ -200,7 +201,8 @@ def query_sites( process_query_results( db, all_files, console, manager, format, page_size, output_file=output_file, job_name=job_name, - auto_confirm=auto_confirm, verbose=verbose + auto_confirm=auto_confirm, verbose=verbose, + async_mode=async_mode ) return CommandResult(success=True, object=all_files, msg=f"Shown {len(all_files)}") @@ -404,7 +406,7 @@ def slice_query( if not collection_name or collection_name == "main": return CommandResult.error("Collection name must be specified and cannot be 'main'") - md_file_pattern = paths + "metadata*.parquet" + md_file_pattern = paths + "*.parquet" try: all_ds_md_files = get_all_files( @@ -519,7 +521,7 @@ def slice_query( try: inferred = infer_metadata_from_files( - owi.local.files_details(target_ds, files_glob="**/metadata_*.parquet", count_rows=True) + owi.local.files_details(target_ds, files_glob="**/*.parquet", count_rows=True) ) for k, v in inferred.items(): target_ds.metadata[k] = v @@ -531,7 +533,17 @@ def slice_query( return CommandResult( success=success, - object={"dataset": target_ds}, + object={ + "dataset": { + "id": target_ds.metadata.get("id"), + "path": target_ds.path, + "collectionName": target_ds.metadata.get("collectionName"), + "access": target_ds.access, + "title": target_ds.metadata.get("title"), + "files": target_ds.metadata.get("fileCount"), + "size": target_ds.metadata.get("totalSize") + } + }, msg=f"Slice {'completed' if success else 'failed'}" ) @@ -614,6 +626,7 @@ def query_aggregate( prefetch: int = 1, page_size: int = 10, output_file: Optional[str] = None, + async_mode: bool = False, console: Optional[Console] = None ) -> CommandResult: """Run SQL aggregation queries across datasets.""" @@ -653,22 +666,50 @@ def query_aggregate( _aggregates = [] processed_count = 0 + total_files = sum([len(v) for v in all_files.values()]) + _files_seen = set() with Progress( SpinnerColumn(), TextColumn("[progress.description]{task.description}"), - TextColumn("[blue]{task.fields[records]} records aggregated"), + "[progress.percentage]{task.percentage:>3.0f}%", + TextColumn("[cyan]{task.fields[files_done]}/{task.fields[files_total]} files"), + TextColumn("[blue]{task.fields[records]:,} records aggregated"), transient=True, console=console ) as progress: - spinner_task = progress.add_task("Aggregating records...", records=processed_count) + spinner_task = progress.add_task( + "Aggregating records...", + total=total_files, + records=processed_count, + files_done=0, + files_total=total_files + ) - for results in db.query_aggregator(): + # Use async or sync iterator based on mode + if async_mode: + from owilix.core.tasks.query_utils import create_async_iterator + results_iterator = create_async_iterator(db) + else: + results_iterator = db.query_aggregator() + + for results in results_iterator: if not results.success: console.print(f"[red]Error when running query: [/red]" + str(results.error)) continue processed_count += len(results.rows) - progress.update(spinner_task, records=processed_count) + + # Track unique files processed + new_files = set(f[0] for f in results.parquet_batch.files) - _files_seen + _files_seen.update(new_files) + files_processed = len(_files_seen) + + progress.update( + spinner_task, + completed=files_processed, + records=processed_count, + files_done=files_processed + ) _aggregates.extend(results.rows) if not _aggregates: @@ -677,6 +718,60 @@ def query_aggregate( aggregates = pd.DataFrame(_aggregates) + # Auto-generate proper re-aggregation query if using default and groupby is specified + final_query = aggregate_query + # Check if using a simple default that doesn't re-aggregate (no GROUP BY in the final query) + needs_reaggregation = ( + groupby and + "group by" not in aggregate_query.lower() and + ("select *" in aggregate_query.lower() or "select * from aggregates" in aggregate_query.lower()) + ) + if needs_reaggregation: + # Parse the select to find count/sum columns that need re-aggregation + # For columns like "count(*) as c", we need to SUM them in the final query + select_parts = [s.strip() for s in select.split(",")] + final_select_parts = [] + + for part in select_parts: + part_lower = part.lower() + # Check if this is an aggregate function + if "count(" in part_lower or "sum(" in part_lower: + # Extract alias if present (e.g., "count(*) as c" -> "c") + if " as " in part_lower: + alias = part.split(" as ")[-1].strip() + # Re-aggregate by summing the partial counts + final_select_parts.append(f"SUM({alias}) as {alias}") + else: + # No alias, use a default + final_select_parts.append(f"SUM({part})") + elif "avg(" in part_lower: + if " as " in part_lower: + alias = part.split(" as ")[-1].strip() + final_select_parts.append(f"AVG({alias}) as {alias}") + else: + final_select_parts.append(part) + elif "min(" in part_lower: + if " as " in part_lower: + alias = part.split(" as ")[-1].strip() + final_select_parts.append(f"MIN({alias}) as {alias}") + else: + final_select_parts.append(part) + elif "max(" in part_lower: + if " as " in part_lower: + alias = part.split(" as ")[-1].strip() + final_select_parts.append(f"MAX({alias}) as {alias}") + else: + final_select_parts.append(part) + else: + # Regular column (likely groupby column), keep as-is + final_select_parts.append(part) + + final_select = ", ".join(final_select_parts) + final_query = f"SELECT {final_select} FROM aggregates GROUP BY {groupby} ORDER BY 2 DESC" + + if format == "table": + console.print(f"[dim]Re-aggregating with: {final_query}[/dim]") + try: if interactive: from IPython import embed @@ -684,7 +779,7 @@ def query_aggregate( "You can run duckdb.query('...from aggregates...')") embed() else: - df = duckdb.query(aggregate_query).to_df() + df = duckdb.query(final_query).to_df() if is_json_mode: json_output = df.to_json(orient="records", lines=True) @@ -694,7 +789,8 @@ def query_aggregate( else: console.print(json_output) else: - console.print(df) + # Show all rows without truncation + console.print(df.to_string()) except Exception as e: logger.exception(f"Error when running query: {e}") @@ -715,6 +811,7 @@ def query_stats( write_back: bool = False, force_write: bool = False, pq_batch_size: int = 10, + async_mode: bool = False, console: Optional[Console] = None ) -> CommandResult: """ @@ -792,6 +889,7 @@ def query_stats( prefetch=1) files_processed = 0 + _files_seen = set() with Progress( SpinnerColumn(), TextColumn("[progress.description]{task.description}"), @@ -809,7 +907,14 @@ def query_stats( files_total=total_files ) - for results in db.query_aggregator(): + # Use async or sync iterator based on mode + if async_mode: + from owilix.core.tasks.query_utils import create_async_iterator + results_iterator = create_async_iterator(db) + else: + results_iterator = db.query_aggregator() + + for results in results_iterator: if not results.success: logger.warning(f"Query error: {results.error}") continue @@ -842,8 +947,10 @@ def query_stats( if max_date is None or date_str > max_date: max_date = date_str - # Update progress with file batch info - files_processed += len(results.parquet_batch.files) + # Track unique files processed (avoid counting same files multiple times) + new_files = set(f[0] for f in results.parquet_batch.files) - _files_seen + _files_seen.update(new_files) + files_processed = len(_files_seen) progress.update( task, completed=files_processed, diff --git a/owilix/core/tasks/query_graphs.py b/owilix/core/tasks/query_graphs.py index 87eac12..f965bc0 100644 --- a/owilix/core/tasks/query_graphs.py +++ b/owilix/core/tasks/query_graphs.py @@ -40,7 +40,7 @@ Common Issues Performance: ------------ +------------ Well equiped Desktop PC: diff --git a/owilix/core/tasks/query_utils.py b/owilix/core/tasks/query_utils.py index 3db84c8..4554113 100644 --- a/owilix/core/tasks/query_utils.py +++ b/owilix/core/tasks/query_utils.py @@ -28,6 +28,64 @@ logger = logging.getLogger(__name__) _DBLOG_CACHE = {} + +def create_async_iterator(db): + """ + Create a streaming iterator from an async executor using queue-based bridge. + + This allows async executors to be consumed synchronously while still + benefiting from async concurrency for the query execution. + + Args: + db: An executor with query_aggregator_async() method + + Returns: + A generator that yields results as they arrive + """ + import asyncio + import queue + import threading + + result_queue = queue.Queue() + error_holder = [None] # Use list to capture exception from thread + + def run_async_collector(): + async def collect_and_queue(): + try: + async for result in db.query_aggregator_async(): + result_queue.put(result) + except Exception as e: + error_holder[0] = e + finally: + result_queue.put(None) # Sentinel to signal completion + try: + await db.close() + except: + pass + + try: + asyncio.run(collect_and_queue()) + except Exception as e: + error_holder[0] = e + result_queue.put(None) # Ensure sentinel is sent on error + + # Start async collection in background thread + collector_thread = threading.Thread(target=run_async_collector, daemon=True) + collector_thread.start() + + # Generator that yields from queue as results arrive + def async_queue_iterator(): + while True: + result = result_queue.get() + if result is None: # Sentinel + # Check if there was an error in the background thread + if error_holder[0] is not None: + raise error_holder[0] + break + yield result + + return async_queue_iterator() + def get_or_create_dblog(job_name: str, logpath: str = './logs') -> Optional[DBLog]: """Get or create DBLog instance for job.""" if not job_name: @@ -449,25 +507,16 @@ def process_query_results( # Open output file if specified output_fh = open(output_file, "w", encoding="utf-8") if output_file and is_json_mode else None - # For async mode, collect all results first then process + # For async mode, use a queue-based bridge to stream results if async_mode: - import asyncio - - async def collect_results(): - results_list = [] - async for result in db.query_aggregator_async(): - results_list.append(result) - await db.close() - return results_list - - all_results = asyncio.run(collect_results()) - results_iterator = iter(all_results) + results_iterator = create_async_iterator(db) else: results_iterator = db.query_aggregator() try: # Start progress tracking show_progress_ui = (format == "table") or (output_file and not with_progress) + if show_progress_ui: progress_display.start(total_files, "Processing files...") diff --git a/owilix/core/tasks/remote.py b/owilix/core/tasks/remote.py index 8712753..86881bc 100644 --- a/owilix/core/tasks/remote.py +++ b/owilix/core/tasks/remote.py @@ -46,7 +46,36 @@ def remote_pull( ignore_data_centers=[push_to_remote] if push_to_remote else None ) - console.print(f"Found {len(datasets)} datasets to pull") + # Print datasets (short format like ls) + for d in datasets: + ds_id = d.metadata.get("internalID", "N/A") or "N/A" + title = d.metadata.get("title", "Untitled") or "Untitled" + + # Dates + start_date = d.metadata.get("startDate", None) + end_date = d.metadata.get("endDate", None) + start = str(start_date)[:10] if start_date else "?" + end = str(end_date)[:10] if end_date else start + date_range = start if start == end else f"{start}→{end}" + + # Location + dc = getattr(d, 'dataCenter', None) or "?" + zone = getattr(d, 'zone', None) or d.metadata.get("zone", "?") or "?" + coll = d.metadata.get("collectionName", "?") or "?" + + # Stats + obj_count = d.metadata.get("objectCount", 0) or 0 + total_size = d.metadata.get("totalSize", 0) or 0 + size_gib = float(total_size) / (1024**3) + access = getattr(d, 'access', None) or d.metadata.get("access", "?") + + line = f"[bold cyan]📦\t{ds_id}[/]\t[bold]{title}[/]\t[dim]{date_range}\t{coll}\t{dc}\t{zone}\t[/]\t#={int(obj_count):,}\t{size_gib:.2f}GiB\t{access}" + console.print(line) + + # Print summary + total_size = sum(d.metadata.get('totalSize', 0) or 0 for d in datasets) + total_files = sum(d.metadata.get('fileCount', 0) or 0 for d in datasets) + console.print(f"\n[bold]Summary:[/bold] {len(datasets)} datasets found, {total_size/1e9:.1f} GB, {total_files:,} files") _add = "" if push_to_remote is None else f" and push to {push_to_remote}" @@ -57,7 +86,7 @@ def remote_pull( files_pattern = [files] for d in datasets: - console.print(f"Fetching files for {d.metadata.title} with file glob filter {files_pattern}") + console.print(f"Fetching files for {d.metadata.get('title', 'Unknown')} with file glob filter {files_pattern}") _remote_files = [rel_can_path(_p, d.path) for _p in manager.remote_data.files(d, files_pattern)] if language is not None: @@ -83,7 +112,7 @@ def remote_pull( _missing_files = set(_remote_files) - set(_local_files) if not overwrite else set(_remote_files) if len(_missing_files) == 0: - console.print(f"Dataset {d.metadata.title} is already up to date. Missing files are {len(_missing_files)}.") + console.print(f"Dataset {d.metadata.get('title', 'Unknown')} is already up to date. Missing files are {len(_missing_files)}.") continue with currentItemProgress() as _progress: @@ -157,9 +186,9 @@ def remote_remove( for d in datasets: try: d.repository.delete(d) - _progress.update(_dtask, advance=1, current_item=f"Deleted {d.metadata.title}") + _progress.update(_dtask, advance=1, current_item=f"Deleted {d.metadata.get('title', 'Unknown')}") except Exception as e: - _progress.update(_error, advance=1, current_item=f"Error deleting {d.metadata.title}: {e.__class__.__name__}: {e}") + _progress.update(_error, advance=1, current_item=f"Error deleting {d.metadata.get('title', 'Unknown')}: {e.__class__.__name__}: {e}") return CommandResult(success=True, object=datasets, msg=f"Removed {len(datasets)} datasets") @@ -168,6 +197,7 @@ def remote_doctor( manager: Any, as_json: bool = False, json_file: Optional[str] = None, + show_tokens: bool = False, console: Optional[Console] = None ) -> CommandResult: """ @@ -242,6 +272,52 @@ def remote_doctor( else: console.print(json.dumps(results)) + # Show token information if requested + if show_tokens: + if not as_json: + console.print("\n[bold]Session Token Information:[/bold]") + try: + if hasattr(manager, 'session') and manager.session: + access_token = manager.session.get_access_token() + refresh_token = manager.session.get_refresh_token() + + if access_token: + # Show truncated token for security (first 20 and last 10 chars) + token_preview = f"{access_token[:20]}...{access_token[-10:]}" if len(access_token) > 30 else access_token + if not as_json: + console.print(f"\tAccess Token: [dim]{token_preview}[/dim]") + console.print(f"\tToken Length: {len(access_token)} chars") + else: + results["access_token"] = {"preview": token_preview, "length": len(access_token)} + else: + if not as_json: + console.print("\t[yellow]Access Token: Not available[/yellow]") + else: + results["access_token"] = None + + if refresh_token: + refresh_preview = f"{refresh_token[:20]}...{refresh_token[-10:]}" if len(refresh_token) > 30 else refresh_token + if not as_json: + console.print(f"\tRefresh Token: [dim]{refresh_preview}[/dim]") + console.print(f"\tRefresh Token Length: {len(refresh_token)} chars") + else: + results["refresh_token"] = {"preview": refresh_preview, "length": len(refresh_token)} + else: + if not as_json: + console.print("\t[yellow]Refresh Token: Not available[/yellow]") + else: + results["refresh_token"] = None + else: + if not as_json: + console.print("\t[yellow]No session available[/yellow]") + else: + results["session"] = None + except Exception as e: + if not as_json: + console.print(f"\t[red]Error getting token info: {e}[/red]") + else: + results["token_error"] = str(e) + return CommandResult(success=True, object=None) @@ -281,7 +357,8 @@ def remote_push( files_pattern = [files] for d in datasets: - console.print(f"Fetching files for {d.title}") + title = d.metadata.get('title', 'Unknown') + console.print(f"Fetching files for {title}") _local_files = [rel_can_path(_p, d.path) for _p in manager.local.files(d, files_pattern)] _ds_remote = manager.remote_data.list(access=d.access, query={"internalID": d.internalID, @@ -296,7 +373,7 @@ def remote_push( repo_names = list(manager.remote_data.get_repos().keys()) raise ValueError(f"Data center {dataCenter} not found. Available: {repo_names}") - console.print(f"Creating dataset {d.title} in datacenter {d.metadata.get('dataCenter')}") + console.print(f"Creating dataset {title} in datacenter {d.metadata.get('dataCenter')}") _ds_remote.append(manager.remote_data.create(d.metadata.get('dataCenter'), **d.metadata)) manager.local.change_id(d, _ds_remote[-1].internalID) @@ -323,11 +400,11 @@ def remote_push( _missing_files = set(_local_files) - set(_remote_files) if not overwrite else set(_local_files) if len(_missing_files) == 0: - console.print(f"Dataset {d.title} is already up to date") + console.print(f"Dataset {d.metadata.get('title', 'Unknown')} is already up to date") continue if dataCenter is not None and _ds_remote[0].dataCenter != dataCenter: - console.print(f"[yellow]Dataset {d.title} is already in datacenter {_ds_remote[0].dataCenter}. Requested {dataCenter} ignored.[/yellow]") + console.print(f"[yellow]Dataset {d.metadata.get('title', 'Unknown')} is already in datacenter {_ds_remote[0].dataCenter}. Requested {dataCenter} ignored.[/yellow]") console.print(f"Uploading {len(_missing_files)} missing files now.") @@ -440,15 +517,15 @@ def remote_diff( console.print("Datasets available only locally:") # Display simple list or title for d in [_d for _d in _local_datasets if get_id(_d) not in _ids]: - console.print(f" - {d.title} ({get_id(d)})") + console.print(f" - {d.metadata.get('title', 'Untitled')} ({get_id(d)})") if len(_ids) < len(_remote_datasets): console.print("Datasets available only remotely:") for d in [_d for _d in _remote_datasets if get_id(_d) not in _ids]: - console.print(f" - {d.title} ({get_id(d)})") + console.print(f" - {d.metadata.get('title', 'Untitled')} ({get_id(d)})") if files is None: - console.print(f"Diff done (for diff on a file level specify files=**/*") + console.print(f"Diff done (for diff on a file level specify --files **/*") return CommandResult(success=True, object=None) try: @@ -465,7 +542,7 @@ def remote_diff( _localds = locals_map[_id] _remoteds = remotes_map[_id] - console.print(f"Metadata diff for dataset {_remoteds.title} (file-diff in calculation)") + console.print(f"Metadata diff for dataset {_remoteds.metadata.get('title', 'Untitled')} (file-diff in calculation)") _keys = set(_localds.metadata.keys()).union(_remoteds.metadata.keys()) _diff_data = [] @@ -477,7 +554,7 @@ def remote_diff( _diff_data.append([key_display, str(l_val), str(r_val)]) # Show table - table = Table(title=f"Metadata Diff: {_localds.title}") + table = Table(title=f"Metadata Diff: {_localds.metadata.get('title', 'Untitled')}") table.add_column("Key") table.add_column("Local") table.add_column("Remote") diff --git a/owilix/core/tasks/warc/__init__.py b/owilix/core/tasks/warc/__init__.py index 65b7596..6b02be4 100644 --- a/owilix/core/tasks/warc/__init__.py +++ b/owilix/core/tasks/warc/__init__.py @@ -12,10 +12,11 @@ The module implements a ZeroMQ PUSH/PULL architecture with telemetry and per-dat ZMQ Sockets ----------- -``` -- ``inproc://records``: PUSH/PULL for WARCTask distribution with back-pressure -- ``inproc://jobs``: PUSH/PULL for FileJob work queue with automatic throttling -- ``inproc://stats``: PUB/SUB for best-effort telemetry and UI updates +:: + +- inproc://records: PUSH/PULL for WARCTask distribution with back-pressure +- inproc://jobs: PUSH/PULL for FileJob work queue with automatic throttling +- inproc://stats: PUB/SUB for best-effort telemetry and UI updates Thread Architecture ------------------- @@ -41,7 +42,7 @@ Key Features CLI Usage Examples -***************** +****************** .. code-block:: @@ -93,10 +94,10 @@ Per-Worker Performance Monitoring ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ - 📊 **Individual Thread Statistics**: Detailed performance metrics for each worker thread - ⚡ **Timing Breakdowns**: File open, seek, read, and write time analysis with standard deviations -- 📡 **: Bandwidth Tracking**: Proper read/write/combined bandwidth calculation (MiB/s) +- 📡 **Bandwidth Tracking**: Proper read/write/combined bandwidth calculation (MiB/s) - 📈 **Activity Monitoring**: Real-time worker utilization and throughput tracking - 🔄 **Performance Variance**: Identification of performance outliers and bottlenecks -- 📊 **: Bandwidth Consistency**: Proper standard deviation analysis for bandwidth stability +- 📊 **Bandwidth Consistency**: Proper standard deviation analysis for bandwidth stability Enhanced Per-Datacenter Analysis with Expanded Timing & Bandwidth ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ @@ -122,8 +123,6 @@ Performance Benefits .. moduleauthor:: Michael Granitzer with support from Claude Sonnet 4.0 .. versionadded:: 1.0 -``` - """ from .query_warc import warc from .parquet_logger import analyze_warc_log \ No newline at end of file diff --git a/owilix/core/tasks/warc/parquet_logger.py b/owilix/core/tasks/warc/parquet_logger.py index 670fc25..fa14f63 100644 --- a/owilix/core/tasks/warc/parquet_logger.py +++ b/owilix/core/tasks/warc/parquet_logger.py @@ -141,14 +141,15 @@ class ParquetJobLogger(object): - **Flexible Storage**: Works with any fsspec-compatible filesystem Directory Structure: - ``` - log_directory/ - ├── process___.parquet - ├── process___.parquet - ├── legacy_job_log.parquet # migrated single files - └── .metadata/ - └── cleanup_info.json - ``` + + :: + + log_directory/ + ├── process___.parquet + ├── process___.parquet + ├── legacy_job_log.parquet # migrated single files + └── .metadata/ + └── cleanup_info.json Performance Characteristics: - get_completion_stats(): O(1) with periodic O(n) refresh @@ -157,20 +158,21 @@ class ParquetJobLogger(object): - Refresh frequency: Configurable (default 60s) Example Usage: - ```python - # Initialize with directory-based logging - logger = DistributedParquetJobLogger( - fs=filesystem, - destination_path="/shared/logs", - batch_size=1000, - refresh_interval_seconds=60 - ) - - # Use exactly like the original logger - logger.log_job_started(job) - stats = logger.get_completion_stats() # Fast, cached access - logger.log_job_completed(job, result) - ``` + + .. code-block:: python + + # Initialize with directory-based logging + logger = DistributedParquetJobLogger( + fs=filesystem, + destination_path="/shared/logs", + batch_size=1000, + refresh_interval_seconds=60 + ) + + # Use exactly like the original logger + logger.log_job_started(job) + stats = logger.get_completion_stats() # Fast, cached access + logger.log_job_completed(job, result) """ def __init__(self, diff --git a/owilix/core/tasks/warc/query_warc.py b/owilix/core/tasks/warc/query_warc.py index 7b81437..bda2f9c 100644 --- a/owilix/core/tasks/warc/query_warc.py +++ b/owilix/core/tasks/warc/query_warc.py @@ -2142,24 +2142,24 @@ def warc(manager, local_specifier: str, remote_specifier: str, urls_file: str = zmq_hwm: int = 1000, resume: bool = False, stats_interval: float = 8.0, group_name:str ='', console: Console = None): """ - WARC Cache fetch function. for a given selected dataset the query is executed and the WARC files for the found URLs are fetched. - The system utilizes message queues and threading in order to optimize bandwith. However, note that if teh number of - records is too large, fetching full warc files might be better. + WARC Cache fetch function. For a given selected dataset the query is executed and the WARC files for the found URLs are fetched. + The system utilizes message queues and threading in order to optimize bandwidth. However, note that if the number of + records is too large, fetching full WARC files might be better. - For optimal performance we also recommend to query only local datastes (i.e. use owilix remote pull before). + For optimal performance we also recommend to query only local datasets (i.e. use owilix remote pull before). - Accessing the WARC Cache requires proper credentials and configuration of the warc cache. please refer to the documentation. + Accessing the WARC Cache requires proper credentials and configuration of the WARC cache. Please refer to the documentation. - The process displays a sophisticated console ui with detailed statistic. + The process displays a sophisticated console UI with detailed statistics. All fetches are logged in a parquet transaction log, such that operations can be resumed. Usage Example: - ``` - python -m owilix.cli query warc --local "all:2025-06-06#1/collectionName=main" \ - where="ows_genai IS NOT NULL AND ows_genai=TRUE" \ - warc_location_cfg=.env-warc-cfg.json urls_file=/home/mgrani/tmp/oem.csv max_workers=20 record_threshold=20000 zmq_hwm=100000 resume=True group_name=all_250606-2D - ``` + .. code-block:: bash + + python -m owilix.cli query warc --local "all:2025-06-06#1/collectionName=main" \\ + where="ows_genai IS NOT NULL AND ows_genai=TRUE" \\ + warc_location_cfg=.env-warc-cfg.json urls_file=/home/mgrani/tmp/oem.csv max_workers=20 record_threshold=20000 zmq_hwm=100000 resume=True group_name=all_250606-2D Args: local_specifier (str): Local filesystem path pattern for input files @@ -2167,7 +2167,7 @@ def warc(manager, local_specifier: str, remote_specifier: str, urls_file: str = urls_file (str, optional): File containing URLs to filter records where (str, optional): SQL WHERE clause for additional filtering limit (int, optional): Maximum number of records to process - files (str): File pattern or list of file patterns to match input Parquet files. Defaults to "**/*.parquet" + files (str): File pattern or list of file patterns to match input Parquet files. Defaults to "``**/*.parquet``" pq_batch_size (int): Parquet file batch size for SQL queries. Defaults to 1. batch_size (int): SQL query batch size. Defaults to 100. prefetch (int): Number of batches to prefetch. Defaults to 10. diff --git a/owilix/core/utils.py b/owilix/core/utils.py index cfd6622..24da0ef 100644 --- a/owilix/core/utils.py +++ b/owilix/core/utils.py @@ -140,22 +140,49 @@ def fill_file_details(fs, files, root=None, cb=None, do_count:bool=True): _root = os.path.commonprefix(files) if root is None else root _total = len(files) for ix, f in enumerate(files): - _returns.append({ - "info": {"size":fs.size(f)},#fs.info(f), + file_info = { + "info": {"size":fs.size(f)}, "path": f, "relpath": os.path.relpath(f, _root), - }) + } if cb is not None: cb(f,ix, _total) - if f.endswith(".parquet") and do_count: # count the number of objects in the parquet file - _returns[-1]["objects"] = 0 - if f.endswith('.parquet'): - try: - with fs.open(f, 'rb') as f: - table = pq.read_table(f) - _returns[-1]['objects'] = table.num_rows - except Exception as e: - _returns[-1]["objects"] = 0 + + if f.endswith(".parquet") and do_count: + + file_info["objects"] = 0 + try: + with fs.open(f, 'rb') as f_obj: + metadata = pq.read_metadata(f_obj) + file_info['objects'] = metadata.num_rows + + # Extract dates from statistics if available + # Look for 'date' or 'warc_date' + names = metadata.schema.names + date_cols = [c for c in names if c in ('date', 'warc_date')] + min_d, max_d = None, None + + for dc in date_cols: + idx = names.index(dc) + for rg in range(metadata.num_row_groups): + col = metadata.row_group(rg).column(idx) + if col.statistics and col.statistics.has_min_max: + _min = col.statistics.min + _max = col.statistics.max + # Decode bytes if needed (pyarrow usually returns appropriate type) + if isinstance(_min, bytes): _min = _min.decode('utf-8') + if isinstance(_max, bytes): _max = _max.decode('utf-8') + + if min_d is None or _min < min_d: min_d = _min + if max_d is None or _max > max_d: max_d = _max + + if min_d: file_info["info"]["startDate"] = min_d + if max_d: file_info["info"]["endDate"] = max_d + + except Exception as e: + file_info["objects"] = 0 + + _returns.append(file_info) return _returns diff --git a/tests/owilix/cli/test_query_slice_integration.py b/tests/owilix/cli/test_query_slice_integration.py new file mode 100644 index 0000000..54fcf0d --- /dev/null +++ b/tests/owilix/cli/test_query_slice_integration.py @@ -0,0 +1,137 @@ +import json +import subprocess +import pytest +import os +import time + +class TestQuerySliceIntegration: + """ + End-to-end integration tests for query slice command. + + Workflow: + 1. Slice 'de' suffix from remote dataset to local. + 2. Verify ID and content. + 3. Slice 'com' suffix from remote dataset into SAME local dataset. + 4. Verify updated content (should increase). + """ + + REMOTE_SPEC = "all/id=0350fecc-e58b-11f0-a8c9-8ebf6bb2cab9" + COLLECTION = "integration_test_slice" + + def run_owi(self, args): + cmd = ["uv", "run", "owi"] + args + print(f"Running: {' '.join(cmd)}") + return subprocess.run(cmd, capture_output=True, text=True, timeout=300) + + def setup_method(self): + # Force cleanup before test + self.cleanup() + + def teardown_method(self): + # Cleanup after test + self.cleanup() + + def cleanup(self): + # Remove any datasets in the test collection + self.run_owi(["local", "rm", f"all/collectionName={self.COLLECTION}", "--yes"]) + + def test_query_slice_integration(self): + print("\n--- Step 1: Initial Slice (DE) ---") + result = self.run_owi([ + "query", "slice", + "-L", self.REMOTE_SPEC, + "--where", "url_suffix='de'", + "--collection", self.COLLECTION, + "--format", "json", + "--yes" + ]) + + assert result.returncode == 0, f"Slice 1 failed: {result.stderr}" + + try: + data = json.loads(result.stdout) + if isinstance(data, list): data = data[0] + dataset_1 = data["dataset"] + except json.JSONDecodeError: + pytest.fail(f"Failed to parse JSON output: {result.stdout}") + + dataset_id = dataset_1["id"] + path_1 = dataset_1["path"] + access_1 = dataset_1["access"] + + print(f"Created Dataset: {dataset_id}") + assert access_1 == "public", f"Expected access='public', got '{access_1}'" + assert self.COLLECTION in path_1, f"Path does not contain collection: {path_1}" + assert os.path.exists(path_1), f"Dataset path does not exist: {path_1}" + + # Verify file count (assuming > 0 for 'de') + # We can simulate verify by checking files content or trusting 'files' count in json? + # The JSON output files count comes from metadata, which should be updated. + # But 'files' input might be 11 (source files), let's check exact output files. + # We can run `local ls` to verify. + + print("\n--- Step 2: Verify Initial Slice ---") + ls_res = self.run_owi([ + "--format", "json", + "local", "ls", + f"all/id={dataset_id}" + ]) + assert ls_res.returncode == 0, f"LS failed: {ls_res.stderr}" + ls_data = json.loads(ls_res.stdout) + if isinstance(ls_data, list): ls_data = ls_data[0] + + files_count_1 = ls_data.get("fileCount", 0) + object_count_1 = ls_data.get("objectCount", 0) + size_1 = ls_data.get("size", 0) + print(f"Step 1 Metadata: Files={files_count_1}, Objects={object_count_1}, Size={size_1}") + + assert files_count_1 > 0, "fileCount should be > 0" + assert size_1 > 0, "totalSize should be > 0" + # objectCount depends on if rows exist, but if files > 0 usually objects > 0 + # unless files are empty. + + start_date = ls_data.get("startDate") + print(f"StartDate: {start_date}") + # Not asserting date yet, just checking output + + + print("\n--- Step 3: Second Slice (COM) - Append ---") + # Slice 'com' into the EXISTING dataset + res_2 = self.run_owi([ + "query", "slice", + "-L", self.REMOTE_SPEC, + "--where", "url_suffix='com'", + "--collection", self.COLLECTION, + "--id", dataset_id, + "--format", "json", + "--yes" + ]) + + assert res_2.returncode == 0, f"Slice 2 failed: {res_2.stderr}" + + try: + data_2 = json.loads(res_2.stdout) + if isinstance(data_2, list): data_2 = data_2[0] + dataset_2 = data_2["dataset"] + except json.JSONDecodeError: + pytest.fail(f"Failed to parse JSON output 2: {res_2.stdout}") + + assert dataset_2["id"] == dataset_id, f"ID changed! Expected {dataset_id}, got {dataset_2['id']}" + + print("\n--- Step 4: Verify Updated Slice ---") + ls_res_2 = self.run_owi([ + "--format", "json", + "local", "ls", + f"all/id={dataset_id}" + ]) + ls_data_2 = json.loads(ls_res_2.stdout) + if isinstance(ls_data_2, list): ls_data_2 = ls_data_2[0] + + files_count_2 = ls_data_2.get("fileCount", 0) + print(f"Step 2 Files: {files_count_2}") + + # We expect file count to increase or stay same if com is empty (unlikely) + # assert files_count_2 >= files_count_1 + # Ideally > if 'com' has data. + + print(f"Test Completed. {dataset_id} verified.")