diff --git a/CHANGELOG.md b/CHANGELOG.md index ce8f7af..81b43a2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,452 +2,459 @@ ## v0.10.0 (2024-06-21) -### Unknown - -* Merge branch 'main' of https://opencode.it4i.eu/openwebsearcheu-public/owi-cli ([`98a381a`](https://github.com/openwebsearcheu-public/owi-cli/commit/98a381a3d1772064082b2c60735e8d4586db9878)) - -## v0.9.0 (2024-06-21) - ### Chore -* chore: added gitlab as vcs for build process ([`9d9a2b7`](https://github.com/openwebsearcheu-public/owi-cli/commit/9d9a2b7bcb7976c51af901071fe1a83e3f9ae197)) +* chore: added gitlab as vcs for build process ([`9d9a2b7`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/9d9a2b7bcb7976c51af901071fe1a83e3f9ae197)) ### Feature -* feat: stream now writes to socket ([`31a2c16`](https://github.com/openwebsearcheu-public/owi-cli/commit/31a2c16974c2e3deacb3233dd657aca0c1e7e427)) +* feat: stream now writes to socket ([`31a2c16`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/31a2c16974c2e3deacb3233dd657aca0c1e7e427)) -* feat: added query stream command - owilix now provides a stream of pyarrow instances from duckdb results ([`1921086`](https://github.com/openwebsearcheu-public/owi-cli/commit/19210866b14c41cadf0f2ea45e532849b408f241)) +* feat: added query stream command - owilix now provides a stream of pyarrow instances from duckdb results ([`1921086`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/19210866b14c41cadf0f2ea45e532849b408f241)) -* feat: added query stream command - owilix now provides a stream of pyarrow instances from duckdb results ([`4afd05f`](https://github.com/openwebsearcheu-public/owi-cli/commit/4afd05ff1b34254c4b54f3de314eabf8162fbd85)) +* feat: added query stream command - owilix now provides a stream of pyarrow instances from duckdb results ([`4afd05f`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/4afd05ff1b34254c4b54f3de314eabf8162fbd85)) -* feat: slice now allows to add to existing dataset and considers proveancne syncing ([`02e24fe`](https://github.com/openwebsearcheu-public/owi-cli/commit/02e24febf5951f308faae9afde373f7ddc78a25d)) +* feat: slice now allows to add to existing dataset and considers proveancne syncing ([`02e24fe`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/02e24febf5951f308faae9afde373f7ddc78a25d)) ### Fix -* fix: corrected merge problem ([`8c35a11`](https://github.com/openwebsearcheu-public/owi-cli/commit/8c35a11501567fcbe3067eff76ea661359249a3f)) +* fix: corrected merge problem ([`8c35a11`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/8c35a11501567fcbe3067eff76ea661359249a3f)) -* fix: retry count for duckdb arrow command ([`e9c4a25`](https://github.com/openwebsearcheu-public/owi-cli/commit/e9c4a255b97dd1a988c3ee540c606b999e91e8de)) +* fix: retry count for duckdb arrow command ([`e9c4a25`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/e9c4a255b97dd1a988c3ee540c606b999e91e8de)) -* fix: added retry count for duckdb query to avoid server errors ([`6b83a2a`](https://github.com/openwebsearcheu-public/owi-cli/commit/6b83a2ad35b12aa0fae1ee5f3ac7d4a134c45367)) +* fix: added retry count for duckdb query to avoid server errors ([`6b83a2a`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/6b83a2ad35b12aa0fae1ee5f3ac7d4a134c45367)) -* fix: added retry count for duckdb query to avoid server errors ([`2e9eb80`](https://github.com/openwebsearcheu-public/owi-cli/commit/2e9eb80c0ef8d02c313221eb04fba3c2ad2bf651)) +* fix: added retry count for duckdb query to avoid server errors ([`2e9eb80`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/2e9eb80c0ef8d02c313221eb04fba3c2ad2bf651)) ### Unknown -* Merge remote-tracking branch 'origin/main' ([`7fc6a50`](https://github.com/openwebsearcheu-public/owi-cli/commit/7fc6a5059e1d765d7100ca25b683e4fb80dcb936)) +* Merge branch 'main' of https://opencode.it4i.eu/openwebsearcheu-public/owi-cli ([`98a381a`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/98a381a3d1772064082b2c60735e8d4586db9878)) -* merge: conflict when arrro duckdb resolved ([`2fa5f62`](https://github.com/openwebsearcheu-public/owi-cli/commit/2fa5f62289744c09502c0771def09d0d6e45359b)) +* Merge remote-tracking branch 'origin/main' ([`7fc6a50`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/7fc6a5059e1d765d7100ca25b683e4fb80dcb936)) -* added retry count ([`4d721c6`](https://github.com/openwebsearcheu-public/owi-cli/commit/4d721c660dbc2c2b62e7b8f8f24b60c7950e610a)) +* merge: conflict when arrro duckdb resolved ([`2fa5f62`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/2fa5f62289744c09502c0771def09d0d6e45359b)) -* Merge remote-tracking branch 'origin/main' ([`59f2ca3`](https://github.com/openwebsearcheu-public/owi-cli/commit/59f2ca383d63b9b1d3df3e503083f526b323df71)) +* added retry count ([`4d721c6`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/4d721c660dbc2c2b62e7b8f8f24b60c7950e610a)) + +* Merge remote-tracking branch 'origin/main' ([`59f2ca3`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/59f2ca383d63b9b1d3df3e503083f526b323df71)) ## v0.8.0 (2024-06-18) ### Chore -* chore: optimized imports ([`89a7c78`](https://github.com/openwebsearcheu-public/owi-cli/commit/89a7c78303d88e9f15c3f04e78eee95fa6c50bcd)) +* chore: optimized imports ([`89a7c78`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/89a7c78303d88e9f15c3f04e78eee95fa6c50bcd)) -* chore: numpy 2.0.0 yielded a problem with libraries in ubuntu ([`1d4ba04`](https://github.com/openwebsearcheu-public/owi-cli/commit/1d4ba043c06c8f873c4af6a9aa09a1f0dd36ae84)) +* chore: numpy 2.0.0 yielded a problem with libraries in ubuntu ([`1d4ba04`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/1d4ba043c06c8f873c4af6a9aa09a1f0dd36ae84)) ### Documentation -* docs: updated doc (a bit) ([`0fe822c`](https://github.com/openwebsearcheu-public/owi-cli/commit/0fe822cdaaf23aed33c0ba8f4553a944568133dd)) +* docs: updated doc (a bit) ([`0fe822c`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/0fe822cdaaf23aed33c0ba8f4553a944568133dd)) -* docs: improved documentation ([`7becf4a`](https://github.com/openwebsearcheu-public/owi-cli/commit/7becf4a7cc2a9eda707de44ef80887fa455f8919)) +* docs: improved documentation ([`7becf4a`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/7becf4a7cc2a9eda707de44ef80887fa455f8919)) ### Feature -* feat: new command: local free to free space (i.e. remove local datasets) ([`3b75b8a`](https://github.com/openwebsearcheu-public/owi-cli/commit/3b75b8a133f05aedec20510e040b7c510c9e4e28)) +* feat: new command: local free to free space (i.e. remove local datasets) ([`3b75b8a`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/3b75b8a133f05aedec20510e040b7c510c9e4e28)) -* feat: added --default_config flag to not load stored config. ([`ff339bf`](https://github.com/openwebsearcheu-public/owi-cli/commit/ff339bfc7f8d8e6f8d3b46462e923971f2a50bc4)) +* feat: added --default_config flag to not load stored config. ([`ff339bf`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/ff339bfc7f8d8e6f8d3b46462e923971f2a50bc4)) ### Fix -* fix: non-existing collection folder yielded an exception when pulling the first time. ([`3aae39f`](https://github.com/openwebsearcheu-public/owi-cli/commit/3aae39f7c479c677deda9695c1ec5c678b68fecf)) +* fix: non-existing collection folder yielded an exception when pulling the first time. ([`3aae39f`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/3aae39f7c479c677deda9695c1ec5c678b68fecf)) -* fix: irods works now when client_zone needs to be selected. ([`2477de8`](https://github.com/openwebsearcheu-public/owi-cli/commit/2477de80813e7050ddb02ebc4c23bc338b0eaae8)) +* fix: irods works now when client_zone needs to be selected. ([`2477de8`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/2477de80813e7050ddb02ebc4c23bc338b0eaae8)) -* fix: irods works now when client_zone needs to be selected. ([`4b8f706`](https://github.com/openwebsearcheu-public/owi-cli/commit/4b8f706b13d227c012de5f75d736922e4f6f3ccb)) +* fix: irods works now when client_zone needs to be selected. ([`4b8f706`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/4b8f706b13d227c012de5f75d736922e4f6f3ccb)) -* fix: documentation was not shown properly ([`050d5a1`](https://github.com/openwebsearcheu-public/owi-cli/commit/050d5a1b99c0b1e0eb982cd31414afadb067ef30)) +* fix: documentation was not shown properly ([`050d5a1`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/050d5a1b99c0b1e0eb982cd31414afadb067ef30)) -* fix: lexis error log not found message implemented. ([`5891717`](https://github.com/openwebsearcheu-public/owi-cli/commit/589171740853d85cf278fc22ed40adb5b262dc9d)) +* fix: lexis error log not found message implemented. ([`5891717`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/589171740853d85cf278fc22ed40adb5b262dc9d)) ### Unknown -* Merge branch 'main' of https://opencode.it4i.eu/openwebsearcheu-public/owi-cli ([`7b1d576`](https://github.com/openwebsearcheu-public/owi-cli/commit/7b1d5763a0fe93029ddbc59923e5c4762ff9b863)) +* Merge branch 'main' of https://opencode.it4i.eu/openwebsearcheu-public/owi-cli ([`7b1d576`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/7b1d5763a0fe93029ddbc59923e5c4762ff9b863)) -* doc: documentation improved ([`1fd05fc`](https://github.com/openwebsearcheu-public/owi-cli/commit/1fd05fcca8225b01d2fb3c6dededc9f3928d6965)) +* doc: documentation improved ([`1fd05fc`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/1fd05fcca8225b01d2fb3c6dededc9f3928d6965)) -* doc: documentation improved ([`8969e26`](https://github.com/openwebsearcheu-public/owi-cli/commit/8969e262a090aee32986d8ac0af52395bbe777f6)) +* doc: documentation improved ([`8969e26`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/8969e262a090aee32986d8ac0af52395bbe777f6)) ## v0.7.4 (2024-06-17) ### Fix -* fix: query now writes to the right collection ([`e85c893`](https://github.com/openwebsearcheu-public/owi-cli/commit/e85c8930f530ad2b27dd0ec06cb102bb474d9123)) +* fix: query now writes to the right collection ([`e85c893`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/e85c8930f530ad2b27dd0ec06cb102bb474d9123)) ### Unknown -* doc: added todo ([`25a6e96`](https://github.com/openwebsearcheu-public/owi-cli/commit/25a6e961ab8a523e286e72ee7c3b2914409d0be0)) +* doc: added todo ([`25a6e96`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/25a6e961ab8a523e286e72ee7c3b2914409d0be0)) + +* doc: small change in metadata default value ([`cf6875c`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/cf6875cdc99c9c8ce5c0c4366f620359aa2e6c75)) -* doc: small change in metadata default value ([`cf6875c`](https://github.com/openwebsearcheu-public/owi-cli/commit/cf6875cdc99c9c8ce5c0c4366f620359aa2e6c75)) ## v0.7.3 (2024-06-16) ### Fix -* fix: speed improvement for irods using queries and direct session for listing objects, bypasing fsspec_irods ([`2f842a1`](https://github.com/openwebsearcheu-public/owi-cli/commit/2f842a16036b48c9456a5f95706e267f4f1bcd9d)) +* fix: speed improvement for irods using queries and direct session for listing objects, bypasing fsspec_irods ([`2f842a1`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/2f842a16036b48c9456a5f95706e267f4f1bcd9d)) ## v0.7.2 (2024-06-16) ### Fix -* fix: slice exits, when no files are found ([`9083708`](https://github.com/openwebsearcheu-public/owi-cli/commit/9083708687d1fd5c6251cb904700e0fdcc76f73a)) +* fix: slice exits, when no files are found ([`9083708`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/9083708687d1fd5c6251cb904700e0fdcc76f73a)) ## v0.7.1 (2024-06-16) ### Fix -* fix: optimze import; pandas no longer on cli start required ([`9bf446b`](https://github.com/openwebsearcheu-public/owi-cli/commit/9bf446bd4c9bc4bd4ab77b4bd515090d88bba825)) +* fix: optimze import; pandas no longer on cli start required ([`9bf446b`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/9bf446bd4c9bc4bd4ab77b4bd515090d88bba825)) ### Unknown -* version update ([`75986f8`](https://github.com/openwebsearcheu-public/owi-cli/commit/75986f8284ccff5bc15e5ff5622a134009d6547a)) +* version update ([`75986f8`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/75986f8284ccff5bc15e5ff5622a134009d6547a)) -* version update ([`7b05707`](https://github.com/openwebsearcheu-public/owi-cli/commit/7b057070969c1d8d84649c931ed0b45e2057ce23)) +* version update ([`7b05707`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/7b057070969c1d8d84649c931ed0b45e2057ce23)) ## v0.7.0 (2024-06-16) ### Chore -* chore: updated py4lexis to 2.1.3 and duckdb to 1.0.0 ([`02e672c`](https://github.com/openwebsearcheu-public/owi-cli/commit/02e672cf882099006eade0501d0565823bd07344)) +* chore: updated py4lexis to 2.1.3 and duckdb to 1.0.0 ([`02e672c`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/02e672cf882099006eade0501d0565823bd07344)) -* chore: moved build to dev dependency group ([`0734502`](https://github.com/openwebsearcheu-public/owi-cli/commit/07345027ab96040ea8e190b9259cee2973e10a40)) +* chore: moved build to dev dependency group ([`0734502`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/07345027ab96040ea8e190b9259cee2973e10a40)) -* chore: removed python tuple version ([`cf78677`](https://github.com/openwebsearcheu-public/owi-cli/commit/cf7867766f9e43c55f51d51951890277fe04b509)) +* chore: removed python tuple version ([`cf78677`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/cf7867766f9e43c55f51d51951890277fe04b509)) ### Feature -* feat: update_metadata command available (but untested due to rights error) ([`39d0408`](https://github.com/openwebsearcheu-public/owi-cli/commit/39d04080610b07f98d00b6c12239481264f43d67)) +* feat: update_metadata command available (but untested due to rights error) ([`39d0408`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/39d04080610b07f98d00b6c12239481264f43d67)) -* feat: slicing finished - when slicing, provenance chain is established and a new dataset is locally created. ([`006f439`](https://github.com/openwebsearcheu-public/owi-cli/commit/006f439187c10a4239755692e5cf2695d4d55551)) +* feat: slicing finished - when slicing, provenance chain is established and a new dataset is locally created. ([`006f439`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/006f439187c10a4239755692e5cf2695d4d55551)) ### Fix -* fix: local ls returns dataset over all collections per default ([`7d037ab`](https://github.com/openwebsearcheu-public/owi-cli/commit/7d037ab148821aab8dac41b4940e84cc4b57d42d)) +* fix: local ls returns dataset over all collections per default ([`7d037ab`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/7d037ab148821aab8dac41b4940e84cc4b57d42d)) -* fix: configuration now supports query federation via zone_path config. Datasets in lrz become available with that setting ([`d6b5813`](https://github.com/openwebsearcheu-public/owi-cli/commit/d6b5813e9ecf3d8b16939c69119634d4c9c687f0)) +* fix: configuration now supports query federation via zone_path config. Datasets in lrz become available with that setting ([`d6b5813`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/d6b5813e9ecf3d8b16939c69119634d4c9c687f0)) -* fix: metadata is now casted after update to avoid list vs. value problems ([`7d4fcb2`](https://github.com/openwebsearcheu-public/owi-cli/commit/7d4fcb279b72ba2c226d67557718402a70dd206c)) +* fix: metadata is now casted after update to avoid list vs. value problems ([`7d4fcb2`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/7d4fcb279b72ba2c226d67557718402a70dd206c)) ### Refactor -* refactor: removed slice command and field ([`57e5637`](https://github.com/openwebsearcheu-public/owi-cli/commit/57e5637e2de8b3179411f9820a127b6afd840ead)) +* refactor: removed slice command and field ([`57e5637`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/57e5637e2de8b3179411f9820a127b6afd840ead)) ### Unknown -* doc: metadata todos ([`8026ef6`](https://github.com/openwebsearcheu-public/owi-cli/commit/8026ef6e5dbcfcf3fb9255b829ad6b2a84996229)) +* doc: metadata todos ([`8026ef6`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/8026ef6e5dbcfcf3fb9255b829ad6b2a84996229)) -* add: summary for dataset table ([`f8faf74`](https://github.com/openwebsearcheu-public/owi-cli/commit/f8faf744d94a1729badd1d321ce2715c53c09e0c)) +* add: summary for dataset table ([`f8faf74`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/f8faf744d94a1729badd1d321ce2715c53c09e0c)) -* add: statistic summary table for datastes ([`b6dac15`](https://github.com/openwebsearcheu-public/owi-cli/commit/b6dac1593cd80a87534588b6c8d3aab69b94ea15)) +* add: statistic summary table for datastes ([`b6dac15`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/b6dac1593cd80a87534588b6c8d3aab69b94ea15)) -* doc: added roadmap; minor formatting of readme ([`6c63d9d`](https://github.com/openwebsearcheu-public/owi-cli/commit/6c63d9de6fd0b446006b709e922b0787ceca3606)) +* doc: added roadmap; minor formatting of readme ([`6c63d9d`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/6c63d9de6fd0b446006b709e922b0787ceca3606)) ## v0.6.1 (2024-06-13) ### Fix -* fix: validation error in metadata when pulling datasets@ ([`f0d467d`](https://github.com/openwebsearcheu-public/owi-cli/commit/f0d467d4d0488fe942a47a1dfbad0ef0a3bc1fee)) +* fix: validation error in metadata when pulling datasets@ ([`f0d467d`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/f0d467d4d0488fe942a47a1dfbad0ef0a3bc1fee)) ### Unknown -* merge with slice changes@ ([`94a44c9`](https://github.com/openwebsearcheu-public/owi-cli/commit/94a44c9f108c7ce6178453c0c60a05c9cf2103cd)) +* merge with slice changes@ ([`94a44c9`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/94a44c9f108c7ce6178453c0c60a05c9cf2103cd)) -* small update ([`7469310`](https://github.com/openwebsearcheu-public/owi-cli/commit/7469310965ee3a70c8dcd84efda7289f732d1949)) +* small update ([`7469310`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/7469310965ee3a70c8dcd84efda7289f732d1949)) -* add: metadata got section for default dataset metadata configuration ([`83c82a5`](https://github.com/openwebsearcheu-public/owi-cli/commit/83c82a57386e53bfa2e1def4496cd2f8c96bb03e)) +* add: metadata got section for default dataset metadata configuration ([`83c82a5`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/83c82a57386e53bfa2e1def4496cd2f8c96bb03e)) -* renamed Spark Indexer in relatedSoftware to OWI Indexer ([`5437933`](https://github.com/openwebsearcheu-public/owi-cli/commit/54379334fa695f52bf6c02239b2a0db4b18a9f2e)) +* renamed Spark Indexer in relatedSoftware to OWI Indexer ([`5437933`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/54379334fa695f52bf6c02239b2a0db4b18a9f2e)) -* add: parquet batch when using duckdb, which allows to have parquet specific template parameters. ([`501128d`](https://github.com/openwebsearcheu-public/owi-cli/commit/501128dbc55ec3b9de30ab6ad401cf45a9d76fed)) +* add: parquet batch when using duckdb, which allows to have parquet specific template parameters. ([`501128d`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/501128dbc55ec3b9de30ab6ad401cf45a9d76fed)) ## v0.6.0 (2024-06-11) ### Feature -* feat: insert with sub-path is possible now. Datasets can be now inserted from larger directories (e.g. year=2023/month=12/day=X) ([`5928228`](https://github.com/openwebsearcheu-public/owi-cli/commit/59282285c9836396feba190f0fa4c278129e89dc)) - -### Fix - -* fix: fixing attempt for semantic release ([`d27a930`](https://github.com/openwebsearcheu-public/owi-cli/commit/d27a930e64c734da85c1c34de9952795827635a1)) +* feat: insert with sub-path is possible now. Datasets can be now inserted from larger directories (e.g. year=2023/month=12/day=X) ([`5928228`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/59282285c9836396feba190f0fa4c278129e89dc)) -* fix: problem with changelog when inserting ([`be0c1d1`](https://github.com/openwebsearcheu-public/owi-cli/commit/be0c1d1ccd09a72932527d79864a0b6b12690830)) +* feat: first version for multithread slicing via duckdb into a single file ([`2221b29`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/2221b299eb623ae56cf8646fa2b3382d7807cae7)) -## v0.5.0 (2024-06-09) +### Fix -### Feature +* fix: fixing attempt for semantic release ([`d27a930`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/d27a930e64c734da85c1c34de9952795827635a1)) -* feat: first version for multithread slicing via duckdb into a single file ([`2221b29`](https://github.com/openwebsearcheu-public/owi-cli/commit/2221b299eb623ae56cf8646fa2b3382d7807cae7)) +* fix: problem with changelog when inserting ([`be0c1d1`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/be0c1d1ccd09a72932527d79864a0b6b12690830)) ## v0.4.0 (2024-06-09) ### Documentation -* docs: readme improved on semantic versioning ([`108fb47`](https://github.com/openwebsearcheu-public/owi-cli/commit/108fb4759e46af2349880286f85a2046846b9b99)) +* docs: readme improved on semantic versioning ([`108fb47`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/108fb4759e46af2349880286f85a2046846b9b99)) ### Feature -* feat: Less command for querying via duckdb ([`086b444`](https://github.com/openwebsearcheu-public/owi-cli/commit/086b444a449c217634511f483baf225cb8a0b2fd)) +* feat: Less command for querying via duckdb ([`086b444`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/086b444a449c217634511f483baf225cb8a0b2fd)) ### Unknown -* doc: documentatoin of duckdb functions improved ([`3369a6b`](https://github.com/openwebsearcheu-public/owi-cli/commit/3369a6b50971ef34b071defc9eada5358713279c)) +* doc: documentatoin of duckdb functions improved ([`3369a6b`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/3369a6b50971ef34b071defc9eada5358713279c)) ## v0.3.0 (2024-06-07) ### Feature -* feat: semantic versioning tested. used this commit to bump ([`eea6e55`](https://github.com/openwebsearcheu-public/owi-cli/commit/eea6e5524eecf369efc8f5f1a0061a8af8cfb7e6)) +* feat: semantic versioning tested. used this commit to bump ([`eea6e55`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/eea6e5524eecf369efc8f5f1a0061a8af8cfb7e6)) ## v0.2.0 (2024-06-07) ### Feature -* feat: semantic versioning ([`ebac113`](https://github.com/openwebsearcheu-public/owi-cli/commit/ebac113fadfa71e750e29ac0d678a21b39595cc1)) +* feat: semantic versioning ([`ebac113`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/ebac113fadfa71e750e29ac0d678a21b39595cc1)) ### Unknown -* gitigore change ([`f9ee469`](https://github.com/openwebsearcheu-public/owi-cli/commit/f9ee46934929bc2bc724e36211a5526bbc17f6f1)) +* gitigore change ([`f9ee469`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/f9ee46934929bc2bc724e36211a5526bbc17f6f1)) ## v0.1.0 (2024-06-07) ### Unknown -* readme example had an error@ ([`7d13a10`](https://github.com/openwebsearcheu-public/owi-cli/commit/7d13a10f458ff492e1940c1a15fc9cf6c45448dc)) +* readme example had an error@ ([`7d13a10`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/7d13a10f458ff492e1940c1a15fc9cf6c45448dc)) -* post build changes ([`5601e9f`](https://github.com/openwebsearcheu-public/owi-cli/commit/5601e9fca586a4dee462138339d503471ea62cd6)) +* post build changes ([`5601e9f`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/5601e9fca586a4dee462138339d503471ea62cd6)) -* poetry problems ([`d9ad3b0`](https://github.com/openwebsearcheu-public/owi-cli/commit/d9ad3b0c4edff66aa5047018a9f701d32efcca94)) +* poetry problems ([`d9ad3b0`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/d9ad3b0c4edff66aa5047018a9f701d32efcca94)) -## v0.5.1 (2024-06-07) +* added config commands and changed metdata field provenance to a list ([`a627ba8`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/a627ba89e155cc426e9c013f290e5f4ff716e107)) -### Chore +* readme updated ([`209ecff`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/209ecff56fad6ed1f2094df66914cb4d646259ab)) -* chore: added upate of versioning in code ([`a97a386`](https://github.com/openwebsearcheu-public/owi-cli/commit/a97a3862096b3e61a03437db9c6b570c11f339a3)) +* fixed error in query not being and for all keys. Missing keys are interpreted as false. ([`51ae9ff`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/51ae9ff5fe062d163cd7ef769b0b37c5857ace50)) -* chore: remove obsolete, old code. cleanup of the module. ([`565bf18`](https://github.com/openwebsearcheu-public/owi-cli/commit/565bf18a4a2f198d10133e136c74644ac27aeeaa)) +* change metadata extraction from title supports multiple title formats. Consequence: more robustness when using lexishttp@ ([`72c7f33`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/72c7f33fe8ae8c039e4ac10cc55cfe2828be830e)) -* chore: added git-changelog; poetry dev and docs groups ([`7a3a076`](https://github.com/openwebsearcheu-public/owi-cli/commit/7a3a076e87ed0c09df35cf11bec9f9cac7448144)) +## v0.5.0 (2024-06-07) -* chore: added git-changelog; poetry dev and docs groups ([`0d550fc`](https://github.com/openwebsearcheu-public/owi-cli/commit/0d550fc1282c42aeca611d928a7a3380f85291df)) +### Unknown -### Feature +* small updates ([`af4b346`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/af4b3468cf9ccd55131144fe3d3bcb947ac1ac03)) -* feat: inserting datasets loally from path ([`a3ef490`](https://github.com/openwebsearcheu-public/owi-cli/commit/a3ef490eaf6826427c852ff8e16f352d9d530ee7)) +* added return value for commands and integrated into logging ([`c4c50c1`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/c4c50c16ea8534559b03f724fb1c5008104c4b51)) -* feat: loca/remote diffs between datasets and files ([`8f9af3d`](https://github.com/openwebsearcheu-public/owi-cli/commit/8f9af3d66e0f081fa2df9a5e789545c28f448b01)) +* added lastChanged metadata and provenance became a list ([`3cab29d`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/3cab29d7e79a7eef26d54509641ee5309f2b6716)) -* feat: added regexp based query in specifier. Field needs to end with * ([`e0da862`](https://github.com/openwebsearcheu-public/owi-cli/commit/e0da862019593b3853b2035a607b1b865f75c058)) +* changed license to Apache 2.0 ([`d8bf773`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/d8bf773b0d2346ca158957aaa3b99fc134fb4331)) -* feat: improved control over displayed columns and sort order ([`b90ff3f`](https://github.com/openwebsearcheu-public/owi-cli/commit/b90ff3f5999cd20c8a4b04a7f762ef0067cf4108)) +* added loging of commands and reading out different log files via ([`02d1479`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/02d14796809cafaa93350692a6d55aa0106d59a9)) -### Unknown +* added HTTPLexisRepository as backup in case irods ports are closed. However, only works for pull / ls. ([`99e65ef`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/99e65efafc25d5671d116fa88ab7a5a163a8fe06)) + +* documented advanced usage ([`7e2023f`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/7e2023f2d86271186027fe0ede0254deffec9c4f)) + +* add pyyaml as dependency for configuration file ([`9b32fa9`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/9b32fa98f913d7ebcc54c6e2c8a2d1ad27ab7e3b)) + +* added relatedSoftware and alternateIdentifier to metadata ([`63c551f`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/63c551ffd140d41f7a16c84f8ac09c9a5eef8606)) + +* added configuration of client zones and zones for enabling federated access to irods zones ([`0a50067`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/0a500679950f157211ea3c405a7b5741d084e758)) + +* intern: fixed path based start_transfer ([`e417de6`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/e417de61ede5ad201f7993587f108a784a0560a3)) + +* added --log_level flag. You can now get all Debug information by using `--log_level=DEBUG` ([`2ca2f5a`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/2ca2f5a9c793d43596419a45e0af629300e0b8f3)) -* added config commands and changed metdata field provenance to a list ([`a627ba8`](https://github.com/openwebsearcheu-public/owi-cli/commit/a627ba89e155cc426e9c013f290e5f4ff716e107)) +* added "owilix" logger ([`b918f4f`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/b918f4f8b3b04e6b679d6fad5fa04111949879b1)) -* readme updated ([`209ecff`](https://github.com/openwebsearcheu-public/owi-cli/commit/209ecff56fad6ed1f2094df66914cb4d646259ab)) +* change ! flag --exclude changed to --remotes ([`deeaf68`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/deeaf68dd9928047181dec3214ef75d64b5b3e23)) -* fixed error in query not being and for all keys. Missing keys are interpreted as false. ([`51ae9ff`](https://github.com/openwebsearcheu-public/owi-cli/commit/51ae9ff5fe062d163cd7ef769b0b37c5857ace50)) +* change config being stored in yaml ([`0f55b49`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/0f55b492de9386a13ea3700fd826a68cb3fd97d8)) -* change metadata extraction from title supports multiple title formats. Consequence: more robustness when using lexishttp@ ([`72c7f33`](https://github.com/openwebsearcheu-public/owi-cli/commit/72c7f33fe8ae8c039e4ac10cc55cfe2828be830e)) +* change config being stored in yaml ([`e85bf6b`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/e85bf6b457a6f348a06e85d74de3433bc75393a9)) -* small updates ([`af4b346`](https://github.com/openwebsearcheu-public/owi-cli/commit/af4b3468cf9ccd55131144fe3d3bcb947ac1ac03)) +* Merge branch 'main' of https://opencode.it4i.eu/openwebsearcheu-public/owi-cli ([`f6d4f71`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/f6d4f71460f656a9cdff760bc7cdbb96d8175a4d)) -* added return value for commands and integrated into logging ([`c4c50c1`](https://github.com/openwebsearcheu-public/owi-cli/commit/c4c50c16ea8534559b03f724fb1c5008104c4b51)) +* add doctor command for checking network connnections configured ([`aa0afd4`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/aa0afd45ae9f0e8cf1ba793555cbedd5026824e8)) -* added lastChanged metadata and provenance became a list ([`3cab29d`](https://github.com/openwebsearcheu-public/owi-cli/commit/3cab29d7e79a7eef26d54509641ee5309f2b6716)) +* fix of a bug on date based selection. ([`bc7a1f9`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/bc7a1f993bb88e4d4228d76ffe345395c12f3c51)) -* changed license to Apache 2.0 ([`d8bf773`](https://github.com/openwebsearcheu-public/owi-cli/commit/d8bf773b0d2346ca158957aaa3b99fc134fb4331)) +* fix dataset creation which is public requires list of rightsURI and list of rights ([`3d3107e`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/3d3107e55734b4af788481128a692b10153647ae)) -* added loging of commands and reading out different log files via ([`02d1479`](https://github.com/openwebsearcheu-public/owi-cli/commit/02d14796809cafaa93350692a6d55aa0106d59a9)) +* intermediate fix for rightsURI requirement on public datasets ([`cdb676d`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/cdb676d8a59dfda92bf3af42b32ae28424085644)) -* added HTTPLexisRepository as backup in case irods ports are closed. However, only works for pull / ls. ([`99e65ef`](https://github.com/openwebsearcheu-public/owi-cli/commit/99e65efafc25d5671d116fa88ab7a5a163a8fe06)) +* update for version 0.4.3 ([`ec08289`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/ec08289664d72a5c03cc20bf43a6784f88c1661f)) -* documented advanced usage ([`7e2023f`](https://github.com/openwebsearcheu-public/owi-cli/commit/7e2023f2d86271186027fe0ede0254deffec9c4f)) +## v0.4.3 (2024-06-03) -* add pyyaml as dependency for configuration file ([`9b32fa9`](https://github.com/openwebsearcheu-public/owi-cli/commit/9b32fa98f913d7ebcc54c6e2c8a2d1ad27ab7e3b)) +### Feature -* added relatedSoftware and alternateIdentifier to metadata ([`63c551f`](https://github.com/openwebsearcheu-public/owi-cli/commit/63c551ffd140d41f7a16c84f8ac09c9a5eef8606)) +* feat: inserting datasets loally from path ([`a3ef490`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/a3ef490eaf6826427c852ff8e16f352d9d530ee7)) -* added configuration of client zones and zones for enabling federated access to irods zones ([`0a50067`](https://github.com/openwebsearcheu-public/owi-cli/commit/0a500679950f157211ea3c405a7b5741d084e758)) +### Unknown -* intern: fixed path based start_transfer ([`e417de6`](https://github.com/openwebsearcheu-public/owi-cli/commit/e417de61ede5ad201f7993587f108a784a0560a3)) +* fix small bug in pushing to remote ([`78114bb`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/78114bb026a759308bd7d78756da85432c345819)) -* added --log_level flag. You can now get all Debug information by using `--log_level=DEBUG` ([`2ca2f5a`](https://github.com/openwebsearcheu-public/owi-cli/commit/2ca2f5a9c793d43596419a45e0af629300e0b8f3)) +## v0.4.2 (2024-06-01) -* added "owilix" logger ([`b918f4f`](https://github.com/openwebsearcheu-public/owi-cli/commit/b918f4f8b3b04e6b679d6fad5fa04111949879b1)) +### Chore -* change ! flag --exclude changed to --remotes ([`deeaf68`](https://github.com/openwebsearcheu-public/owi-cli/commit/deeaf68dd9928047181dec3214ef75d64b5b3e23)) +* chore: added upate of versioning in code ([`a97a386`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/a97a3862096b3e61a03437db9c6b570c11f339a3)) -* change config being stored in yaml ([`0f55b49`](https://github.com/openwebsearcheu-public/owi-cli/commit/0f55b492de9386a13ea3700fd826a68cb3fd97d8)) +* chore: remove obsolete, old code. cleanup of the module. ([`565bf18`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/565bf18a4a2f198d10133e136c74644ac27aeeaa)) -* change config being stored in yaml ([`e85bf6b`](https://github.com/openwebsearcheu-public/owi-cli/commit/e85bf6b457a6f348a06e85d74de3433bc75393a9)) +* chore: added git-changelog; poetry dev and docs groups ([`7a3a076`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/7a3a076e87ed0c09df35cf11bec9f9cac7448144)) -* Merge branch 'main' of https://opencode.it4i.eu/openwebsearcheu-public/owi-cli ([`f6d4f71`](https://github.com/openwebsearcheu-public/owi-cli/commit/f6d4f71460f656a9cdff760bc7cdbb96d8175a4d)) +* chore: added git-changelog; poetry dev and docs groups ([`0d550fc`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/0d550fc1282c42aeca611d928a7a3380f85291df)) -* add doctor command for checking network connnections configured ([`aa0afd4`](https://github.com/openwebsearcheu-public/owi-cli/commit/aa0afd45ae9f0e8cf1ba793555cbedd5026824e8)) +### Feature -* fix of a bug on date based selection. ([`bc7a1f9`](https://github.com/openwebsearcheu-public/owi-cli/commit/bc7a1f993bb88e4d4228d76ffe345395c12f3c51)) +* feat: loca/remote diffs between datasets and files ([`8f9af3d`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/8f9af3d66e0f081fa2df9a5e789545c28f448b01)) -* fix dataset creation which is public requires list of rightsURI and list of rights ([`3d3107e`](https://github.com/openwebsearcheu-public/owi-cli/commit/3d3107e55734b4af788481128a692b10153647ae)) +* feat: added regexp based query in specifier. Field needs to end with * ([`e0da862`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/e0da862019593b3853b2035a607b1b865f75c058)) -* intermediate fix for rightsURI requirement on public datasets ([`cdb676d`](https://github.com/openwebsearcheu-public/owi-cli/commit/cdb676d8a59dfda92bf3af42b32ae28424085644)) +* feat: improved control over displayed columns and sort order ([`b90ff3f`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/b90ff3f5999cd20c8a4b04a7f762ef0067cf4108)) -* update for version 0.4.3 ([`ec08289`](https://github.com/openwebsearcheu-public/owi-cli/commit/ec08289664d72a5c03cc20bf43a6784f88c1661f)) +### Unknown + +* added versionoing ([`63377ff`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/63377ff0b49b5239f3bf7b8b04e6078aaf356183)) -* fix small bug in pushing to remote ([`78114bb`](https://github.com/openwebsearcheu-public/owi-cli/commit/78114bb026a759308bd7d78756da85432c345819)) +## v0.4.1 (2024-05-31) -* added versionoing ([`63377ff`](https://github.com/openwebsearcheu-public/owi-cli/commit/63377ff0b49b5239f3bf7b8b04e6078aaf356183)) +### Unknown -* compatibility with py4lexis-2.1.1 ([`79ef75e`](https://github.com/openwebsearcheu-public/owi-cli/commit/79ef75e954491b412479de5e083eebf8a79e51fb)) +* compatibility with py4lexis-2.1.1 ([`79ef75e`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/79ef75e954491b412479de5e083eebf8a79e51fb)) -* readme update for admin command ([`b231051`](https://github.com/openwebsearcheu-public/owi-cli/commit/b2310515bf59dbe4752d3ff78a1216be313432a2)) +* readme update for admin command ([`b231051`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/b2310515bf59dbe4752d3ff78a1216be313432a2)) -* readme update ([`de2c6ed`](https://github.com/openwebsearcheu-public/owi-cli/commit/de2c6eddef8949bb3e4c338dde18c01ecef64f7c)) +* readme update ([`de2c6ed`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/de2c6eddef8949bb3e4c338dde18c01ecef64f7c)) -* bad merge ([`6c71f56`](https://github.com/openwebsearcheu-public/owi-cli/commit/6c71f569de4b21f5600af4d6932340c7b456c1fd)) +* bad merge ([`6c71f56`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/6c71f569de4b21f5600af4d6932340c7b456c1fd)) -* added admin set_irods_metadata command ([`45910b5`](https://github.com/openwebsearcheu-public/owi-cli/commit/45910b5655ae3df3ff313414943df5f1e46d9ee3)) +* added admin set_irods_metadata command ([`45910b5`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/45910b5655ae3df3ff313414943df5f1e46d9ee3)) -* intermediate status for getting all lexis to work ([`3686911`](https://github.com/openwebsearcheu-public/owi-cli/commit/3686911fc50e4a435b8ebb8b758239fce4d375fb)) +* intermediate status for getting all lexis to work ([`3686911`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/3686911fc50e4a435b8ebb8b758239fce4d375fb)) -* dataset creation alpha. API not working ([`6c67198`](https://github.com/openwebsearcheu-public/owi-cli/commit/6c671987e4fd7fd05b17bf5936cbedf8af2d7fc1)) +* dataset creation alpha. API not working ([`6c67198`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/6c671987e4fd7fd05b17bf5936cbedf8af2d7fc1)) -* update local ls; readme ([`d21412b`](https://github.com/openwebsearcheu-public/owi-cli/commit/d21412b99611350093d129d5f146fb662205ddc7)) +* update local ls; readme ([`d21412b`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/d21412b99611350093d129d5f146fb662205ddc7)) -* bit of cleanup ([`743dd22`](https://github.com/openwebsearcheu-public/owi-cli/commit/743dd22002c79425fe4cf7d62c635663e61d26ca)) +* bit of cleanup ([`743dd22`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/743dd22002c79425fe4cf7d62c635663e61d26ca)) -* insert /update of local dataset form path down ([`f29a130`](https://github.com/openwebsearcheu-public/owi-cli/commit/f29a13038fe5b19ec47ea3dac4fb1e1017fa492c)) +* insert /update of local dataset form path down ([`f29a130`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/f29a13038fe5b19ec47ea3dac4fb1e1017fa492c)) -* Merge branch 'develop' of https://opencode.it4i.eu/openwebsearcheu-public/owi-cli into develop ([`00492e2`](https://github.com/openwebsearcheu-public/owi-cli/commit/00492e266a8c0631b4e51ab487e6c14ac9059bdb)) +* Merge branch 'develop' of https://opencode.it4i.eu/openwebsearcheu-public/owi-cli into develop ([`00492e2`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/00492e266a8c0631b4e51ab487e6c14ac9059bdb)) -* swichted to repositories ([`2b933aa`](https://github.com/openwebsearcheu-public/owi-cli/commit/2b933aa5c92bf7c27f132c2a5e00fee428ebbfd0)) +* swichted to repositories ([`2b933aa`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/2b933aa5c92bf7c27f132c2a5e00fee428ebbfd0)) -* small changes. ([`3a10529`](https://github.com/openwebsearcheu-public/owi-cli/commit/3a10529fddfd21bfdcd635c4ee12dc7547bbcfe5)) +* small changes. ([`3a10529`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/3a10529fddfd21bfdcd635c4ee12dc7547bbcfe5)) -* error in install ([`63ee514`](https://github.com/openwebsearcheu-public/owi-cli/commit/63ee5142dd8aed67f13d02efc4a738e192207636)) +* error in install ([`63ee514`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/63ee5142dd8aed67f13d02efc4a738e192207636)) -* Merge branch 'main' into develop ([`c0e28e8`](https://github.com/openwebsearcheu-public/owi-cli/commit/c0e28e8f425d8880f8da618914164a31312ac468)) +* Merge branch 'main' into develop ([`c0e28e8`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/c0e28e8f425d8880f8da618914164a31312ac468)) -* docu update ([`b1d28e1`](https://github.com/openwebsearcheu-public/owi-cli/commit/b1d28e128a44ae93b0b906810a2d9a6dff3e22ba)) +* docu update ([`b1d28e1`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/b1d28e128a44ae93b0b906810a2d9a6dff3e22ba)) -* udpate readme install ([`df9f687`](https://github.com/openwebsearcheu-public/owi-cli/commit/df9f687f23934da24eec6945c31b97fd5c9bcb35)) +* udpate readme install ([`df9f687`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/df9f687f23934da24eec6945c31b97fd5c9bcb35)) -* added prepration flag ([`82b5171`](https://github.com/openwebsearcheu-public/owi-cli/commit/82b5171d619ed433f5d274329a89e15a838c8abb)) +* added prepration flag ([`82b5171`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/82b5171d619ed433f5d274329a89e15a838c8abb)) -* added bug in pulling datasets ([`86e10a0`](https://github.com/openwebsearcheu-public/owi-cli/commit/86e10a038c43e6dd6f959aecaad37a0fe3a6abf7)) +* added bug in pulling datasets ([`86e10a0`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/86e10a038c43e6dd6f959aecaad37a0fe3a6abf7)) -* added dataset creation and metadta admin panel which metadata chekc. Bot incomplete, since server api is not working properly ([`e559d6b`](https://github.com/openwebsearcheu-public/owi-cli/commit/e559d6bfdcee85d8dbda93d22f5acc14f8ecdedf)) +* added dataset creation and metadta admin panel which metadata chekc. Bot incomplete, since server api is not working properly ([`e559d6b`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/e559d6bfdcee85d8dbda93d22f5acc14f8ecdedf)) -* added metadata management; added create v0.01 ([`40912de`](https://github.com/openwebsearcheu-public/owi-cli/commit/40912de7a5367a94e35b0fba623e37f78bb7272c)) +* added metadata management; added create v0.01 ([`40912de`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/40912de7a5367a94e35b0fba623e37f78bb7272c)) -* reworked commands. Converted to recent metadta format ([`e8a35e0`](https://github.com/openwebsearcheu-public/owi-cli/commit/e8a35e0cdf97b9fb5a039ba056d4443fac8734eb)) +* reworked commands. Converted to recent metadta format ([`e8a35e0`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/e8a35e0cdf97b9fb5a039ba056d4443fac8734eb)) -* reworked commands. Converted to recent metadta format ([`bf55504`](https://github.com/openwebsearcheu-public/owi-cli/commit/bf55504202d245073b9935381e5f92a4668f2522)) +* reworked commands. Converted to recent metadta format ([`bf55504`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/bf55504202d245073b9935381e5f92a4668f2522)) -* removed path from stats ([`cc400d6`](https://github.com/openwebsearcheu-public/owi-cli/commit/cc400d6891d3ec346b62ad23ab4ee52cab45576e)) +* removed path from stats ([`cc400d6`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/cc400d6891d3ec346b62ad23ab4ee52cab45576e)) -* added todo ([`b4da4fa`](https://github.com/openwebsearcheu-public/owi-cli/commit/b4da4fa43cea51622b693b660790373c00ca0530)) +* added todo ([`b4da4fa`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/b4da4fa43cea51622b693b660790373c00ca0530)) -* docu;improved get_dataset; irods imp; test; fillet cmd; ([`f32db96`](https://github.com/openwebsearcheu-public/owi-cli/commit/f32db965d570341d932257716a843858cfb1513e)) +* docu;improved get_dataset; irods imp; test; fillet cmd; ([`f32db96`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/f32db965d570341d932257716a843858cfb1513e)) -* added documentatoin ([`76df3e9`](https://github.com/openwebsearcheu-public/owi-cli/commit/76df3e9ab33fb171d638cb4226f5d1c9e123b9ed)) +* added documentatoin ([`76df3e9`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/76df3e9ab33fb171d638cb4226f5d1c9e123b9ed)) -* update commands + documentaiton; update dependencies with +rich -tqdm ([`41ccc99`](https://github.com/openwebsearcheu-public/owi-cli/commit/41ccc9919c4881a856ea60e8cea5cb28558d59e1)) +* update commands + documentaiton; update dependencies with +rich -tqdm ([`41ccc99`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/41ccc9919c4881a856ea60e8cea5cb28558d59e1)) -* added first sql template mechanisms plus scripts ([`dd7df70`](https://github.com/openwebsearcheu-public/owi-cli/commit/dd7df7015a59b0093fa219a0c9b643cf7900fc45)) +* added first sql template mechanisms plus scripts ([`dd7df70`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/dd7df7015a59b0093fa219a0c9b643cf7900fc45)) -* update commands. middle of change ([`af4f5ac`](https://github.com/openwebsearcheu-public/owi-cli/commit/af4f5ac79420f6860d3d9ed0eaad0ac5b6f58776)) +* update commands. middle of change ([`af4f5ac`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/af4f5ac79420f6860d3d9ed0eaad0ac5b6f58776)) -* prpeared pyproject-local-dev.toml to work with py4lexis locally (needed to use python 3.12) ([`667889b`](https://github.com/openwebsearcheu-public/owi-cli/commit/667889b9753f05f142277c49b6dd6765d19ea4c5)) +* prpeared pyproject-local-dev.toml to work with py4lexis locally (needed to use python 3.12) ([`667889b`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/667889b9753f05f142277c49b6dd6765d19ea4c5)) -* updated readme. Added changelog ([`0956f65`](https://github.com/openwebsearcheu-public/owi-cli/commit/0956f655d4a91cc95632cc5465d2a7a19eac4fd0)) +* updated readme. Added changelog ([`0956f65`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/0956f655d4a91cc95632cc5465d2a7a19eac4fd0)) -* typo ([`fb80821`](https://github.com/openwebsearcheu-public/owi-cli/commit/fb8082194a6266e5f6dccb2384b41c9faaf32ab4)) +* typo ([`fb80821`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/fb8082194a6266e5f6dccb2384b41c9faaf32ab4)) -* update poetry config and bump to version ([`d3fcfea`](https://github.com/openwebsearcheu-public/owi-cli/commit/d3fcfea00cbc0ee2589958b84da0b628694178d4)) +* update poetry config and bump to version ([`d3fcfea`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/d3fcfea00cbc0ee2589958b84da0b628694178d4)) -* added duckdb shell ([`e068e76`](https://github.com/openwebsearcheu-public/owi-cli/commit/e068e76fca3503d127b9ce258d5bcaa47138d65c)) +* added duckdb shell ([`e068e76`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/e068e76fca3503d127b9ce258d5bcaa47138d65c)) -* final alpha version for pull ([`b0d6bfd`](https://github.com/openwebsearcheu-public/owi-cli/commit/b0d6bfd6787c3b37251b8ad622bba44a26eebe90)) +* final alpha version for pull ([`b0d6bfd`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/b0d6bfd6787c3b37251b8ad622bba44a26eebe90)) -* first version working ([`e223f92`](https://github.com/openwebsearcheu-public/owi-cli/commit/e223f92c4151fc569d86c9ca6e2487fc412d05dd)) +* first version working ([`e223f92`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/e223f92c4151fc569d86c9ca6e2487fc412d05dd)) -* Threading. in the middle ([`29cab7e`](https://github.com/openwebsearcheu-public/owi-cli/commit/29cab7e9074d0efeb7cb6dc686e6c19348b0d453)) +* Threading. in the middle ([`29cab7e`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/29cab7e9074d0efeb7cb6dc686e6c19348b0d453)) -* added file based transfer ([`199bded`](https://github.com/openwebsearcheu-public/owi-cli/commit/199bded0599c6503e0db6052e1c3bcc72780f96f)) +* added file based transfer ([`199bded`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/199bded0599c6503e0db6052e1c3bcc72780f96f)) -* added duckdb plugin (first test); changed download to thread based download ([`ad7d3bc`](https://github.com/openwebsearcheu-public/owi-cli/commit/ad7d3bc26f1000dc7ec77c934f25f55fbc1d36ec)) +* added duckdb plugin (first test); changed download to thread based download ([`ad7d3bc`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/ad7d3bc26f1000dc7ec77c934f25f55fbc1d36ec)) -* swichted to click and rich ([`ac16ea5`](https://github.com/openwebsearcheu-public/owi-cli/commit/ac16ea558fe9b60eb114df7a6908885c343f0ab1)) +* swichted to click and rich ([`ac16ea5`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/ac16ea558fe9b60eb114df7a6908885c343f0ab1)) -* added optional duckdb dependency ([`b0590ce`](https://github.com/openwebsearcheu-public/owi-cli/commit/b0590ce5b62d036c8dcae79ffdc0a255a1f0dcd2)) +* added optional duckdb dependency ([`b0590ce`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/b0590ce5b62d036c8dcae79ffdc0a255a1f0dcd2)) -* update readme; some testing with poetry ([`9941e00`](https://github.com/openwebsearcheu-public/owi-cli/commit/9941e0055decb8c6c6791ac602b3f9f648959c42)) +* update readme; some testing with poetry ([`9941e00`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/9941e0055decb8c6c6791ac602b3f9f648959c42)) -* added local commands;swithc to poetry;first build ([`564e437`](https://github.com/openwebsearcheu-public/owi-cli/commit/564e4377c43c05cf2cea1383bb00aaee9fd55a49)) +* added local commands;swithc to poetry;first build ([`564e437`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/564e4377c43c05cf2cea1383bb00aaee9fd55a49)) -* added plugin for testing ([`0072b4a`](https://github.com/openwebsearcheu-public/owi-cli/commit/0072b4af71a31d16a17a0db9479b27a36c9db85c)) +* added plugin for testing ([`0072b4a`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/0072b4af71a31d16a17a0db9479b27a36c9db85c)) -* path corrected ([`c0fb2ac`](https://github.com/openwebsearcheu-public/owi-cli/commit/c0fb2acf619475693d22f06792e5768173b9e2ef)) +* path corrected ([`c0fb2ac`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/c0fb2acf619475693d22f06792e5768173b9e2ef)) -* bug fix for tar extraction ([`8290610`](https://github.com/openwebsearcheu-public/owi-cli/commit/82906107aaa556ad381ffe5f32af4807e060882c)) +* bug fix for tar extraction ([`8290610`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/82906107aaa556ad381ffe5f32af4807e060882c)) -* bug-fix: dataset get ([`77a722c`](https://github.com/openwebsearcheu-public/owi-cli/commit/77a722c67d5e146c53f88f660252cb05dcb4965e)) +* bug-fix: dataset get ([`77a722c`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/77a722c67d5e146c53f88f660252cb05dcb4965e)) -* added field list for showing dataset dataframe; first version for extract tar.gz file ([`d001e67`](https://github.com/openwebsearcheu-public/owi-cli/commit/d001e67d0886e8420b4df471ebcfca9f7df74f45)) +* added field list for showing dataset dataframe; first version for extract tar.gz file ([`d001e67`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/d001e67d0886e8420b4df471ebcfca9f7df74f45)) -* added more verbose messages in pull ([`79b07c6`](https://github.com/openwebsearcheu-public/owi-cli/commit/79b07c6a9a88a788176c220966959693655e182a)) +* added more verbose messages in pull ([`79b07c6`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/79b07c6a9a88a788176c220966959693655e182a)) -* added command line option for default path ([`a604506`](https://github.com/openwebsearcheu-public/owi-cli/commit/a60450609d9c15101f1d08a644cc8f139f6639b1)) +* added command line option for default path ([`a604506`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/a60450609d9c15101f1d08a644cc8f139f6639b1)) -* file based filtering (but not download) ([`da59c2b`](https://github.com/openwebsearcheu-public/owi-cli/commit/da59c2b885fe3bd495fdf25dfc687c30325f5247)) +* file based filtering (but not download) ([`da59c2b`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/da59c2b885fe3bd495fdf25dfc687c30325f5247)) -* update setup ([`a23e4d2`](https://github.com/openwebsearcheu-public/owi-cli/commit/a23e4d29b7de97dcb42a77f487600f066bac97ee)) +* update setup ([`a23e4d2`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/a23e4d29b7de97dcb42a77f487600f066bac97ee)) -* setup and command line ([`78e7675`](https://github.com/openwebsearcheu-public/owi-cli/commit/78e767560bd23fadaf281563a453808a39b0d62f)) +* setup and command line ([`78e7675`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/78e767560bd23fadaf281563a453808a39b0d62f)) -* udpated towards specifier ([`01af31e`](https://github.com/openwebsearcheu-public/owi-cli/commit/01af31e0a5c84443556e1b763ae5152893634068)) +* udpated towards specifier ([`01af31e`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/01af31e0a5c84443556e1b763ae5152893634068)) -* rmoeved submodule ([`51f5815`](https://github.com/openwebsearcheu-public/owi-cli/commit/51f581556bf46733c51f14c5bd8b03032ed1a5c8)) +* rmoeved submodule ([`51f5815`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/51f581556bf46733c51f14c5bd8b03032ed1a5c8)) -* rmoeved submodule ([`472b5a1`](https://github.com/openwebsearcheu-public/owi-cli/commit/472b5a15a9216abf3669f5dbbacef4f93d4ac19c)) +* rmoeved submodule ([`472b5a1`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/472b5a15a9216abf3669f5dbbacef4f93d4ac19c)) -* Removed submodule ([`4d2deb1`](https://github.com/openwebsearcheu-public/owi-cli/commit/4d2deb1da4e7aaecf9a250d035f1d88a0e2dd42e)) +* Removed submodule ([`4d2deb1`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/4d2deb1da4e7aaecf9a250d035f1d88a0e2dd42e)) -* first cli version: dataset list, stats ([`e271636`](https://github.com/openwebsearcheu-public/owi-cli/commit/e271636964def01fee8de7f516dc7509f04e77f1)) +* first cli version: dataset list, stats ([`e271636`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/e271636964def01fee8de7f516dc7509f04e77f1)) -* first cli version: dataset list, stats ([`e4bff1c`](https://github.com/openwebsearcheu-public/owi-cli/commit/e4bff1cf7e7f541bccb9be4cce5e1db27430c97b)) +* first cli version: dataset list, stats ([`e4bff1c`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/e4bff1cf7e7f541bccb9be4cce5e1db27430c97b)) -* submodule change ([`5e57239`](https://github.com/openwebsearcheu-public/owi-cli/commit/5e57239e9890fc9c49e632ffb2d4e49b1ef3db87)) +* submodule change ([`5e57239`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/5e57239e9890fc9c49e632ffb2d4e49b1ef3db87)) -* added workflow example from Lexis plus authentication ([`7ce7052`](https://github.com/openwebsearcheu-public/owi-cli/commit/7ce70522179994cf705e848e9d70f9912bf776a0)) +* added workflow example from Lexis plus authentication ([`7ce7052`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/7ce70522179994cf705e848e9d70f9912bf776a0)) -* renamed main to datasets ([`e88f553`](https://github.com/openwebsearcheu-public/owi-cli/commit/e88f5536825dc781f3471163db012d69b6a2aac2)) +* renamed main to datasets ([`e88f553`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/e88f5536825dc781f3471163db012d69b6a2aac2)) -* datasets api working ([`c9f742d`](https://github.com/openwebsearcheu-public/owi-cli/commit/c9f742d241426e3a291649dca30abd6babca43ba)) +* datasets api working ([`c9f742d`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/c9f742d241426e3a291649dca30abd6babca43ba)) -* added openwebsearch/py4lexis as submodule; update readme and test ([`fa36598`](https://github.com/openwebsearcheu-public/owi-cli/commit/fa36598886c5a1ce8e7d2c66ce57196f8de64892)) +* added openwebsearch/py4lexis as submodule; update readme and test ([`fa36598`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/fa36598886c5a1ce8e7d2c66ce57196f8de64892)) -* added username /pwd configuration via Enviornment Variables LEXIS_USERNAME and LEXIS_PASSWORD ([`093ef31`](https://github.com/openwebsearcheu-public/owi-cli/commit/093ef31eac636b235695a4e6365b712995ee3919)) +* added username /pwd configuration via Enviornment Variables LEXIS_USERNAME and LEXIS_PASSWORD ([`093ef31`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/093ef31eac636b235695a4e6365b712995ee3919)) -* removed Py4Lexis in order to fork the repo ([`55a61ef`](https://github.com/openwebsearcheu-public/owi-cli/commit/55a61efc30458fd7901e1c78c78561b95ba4808c)) +* removed Py4Lexis in order to fork the repo ([`55a61ef`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/55a61efc30458fd7901e1c78c78561b95ba4808c)) -* added git modules ([`2f7a8eb`](https://github.com/openwebsearcheu-public/owi-cli/commit/2f7a8eb0764770f18733dc96290f00adc917b3ab)) +* added git modules ([`2f7a8eb`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/2f7a8eb0764770f18733dc96290f00adc917b3ab)) -* added howto for getting account password ([`db0a408`](https://github.com/openwebsearcheu-public/owi-cli/commit/db0a408e00400060045d0d619d7b58d2e4430c78)) +* added howto for getting account password ([`db0a408`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/db0a408e00400060045d0d619d7b58d2e4430c78)) -* Initial commit ([`d06f6f7`](https://github.com/openwebsearcheu-public/owi-cli/commit/d06f6f74a37270dd9d50814783eaf84a97541f4f)) +* Initial commit ([`d06f6f7`](https://opencode.it4i.eu/openwebsearcheu-public/owi-cli/-/commit/d06f6f74a37270dd9d50814783eaf84a97541f4f)) diff --git a/owilix/cmd/query.py b/owilix/cmd/query.py index 5e24ffb..bc58805 100644 --- a/owilix/cmd/query.py +++ b/owilix/cmd/query.py @@ -42,7 +42,7 @@ class QueryCommands(BaseCommand): Fetches all files from the datasets selected by local and remote specifiers and grouped by filesystem. Returns: - Dict[AbstractFileSystem, List[(str,str)]]: A dictionary with filesystems as keys and list of tuples + Dict[AbstractFileSystem, List[(str,str)]], List[Datasets]: A dictionary with filesystems as keys and list of tuples with file and dataset path as values. """ @@ -339,6 +339,44 @@ def stream(self, local_specifier: str, remote_specifier: str, if writer: writer.close() + def configure_kafka_producer(): + from confluent_kafka import Producer + import json + + # Kafka Configuration + KAFKA_BOOTSTRAP_SERVERS = 'localhost:9092' + KAFKA_TOPIC = 'pyarrow_stream' + + # Kafka Producer configuration + producer_config = { + 'bootstrap.servers': KAFKA_BOOTSTRAP_SERVERS, + 'client.id': 'pyarrow-stream-producer' + } + + # Create a Kafka producer + return Producer(producer_config) + + def stream_to_kafka(buffer_queue, topic): + while True: + batch = buffer_queue.get() + if batch is None: + break # End of stream + + # Serialize the batch to a bytes buffer + sink = pa.BufferOutputStream() + writer = ipc.RecordBatchStreamWriter(sink, batch.schema) + writer.write_table(batch) + writer.close() + + # Get the serialized data as bytes + data = sink.getvalue().to_pybytes() + + # Send serialized data to Kafka + producer.produce(topic, data) + producer.flush() # Ensure the message is sent + + buffer_queue.task_done() + all_files = self.get_all_files(files, local_specifier, remote_specifier, print_it=False) _datasets = all_files.keys() all_files = self.group_all_files_by_fs(all_files) diff --git a/pyproject.toml b/pyproject.toml index 5c6d0a0..e24e42f 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -74,11 +74,11 @@ version_pattern = [ ] major_on_zero = false branch = "main" -upload_to_PyPI = true -upload_to_release = true +upload_to_vcs_release = true build_command = "pip install poetry && poetry build" -repository_url = "https://opencode.it4i.eu/openwebsearcheu-public/owi-cli" -vcs_provider = "gitlab" + +[tool.semantic_release.remote] +domain = "https://opencode.it4i.eu" type="gitlab" -- 2.51.2 From 06b6e1e7c34abe51a8aebef1e38aeecb5f084938 Mon Sep 17 00:00:00 2001 From: Michael Granitzer Date: Fri, 21 Jun 2024 15:02:26 +0200 Subject: [PATCH 2/3] fix: added typecasting for parameter calls --- owilix/cmd/base.py | 80 +++++++++++++++++++++++++++++++++++++++-- owilix/cmd/query.py | 4 +-- owilix/core/metadata.py | 1 + pyproject.toml | 1 + 4 files changed, 81 insertions(+), 5 deletions(-) diff --git a/owilix/cmd/base.py b/owilix/cmd/base.py index 52a1f3c..a915032 100644 --- a/owilix/cmd/base.py +++ b/owilix/cmd/base.py @@ -5,7 +5,10 @@ import re from collections import Counter from logging.handlers import RotatingFileHandler from statistics import stdev -from typing import Iterable +from typing import Iterable, Tuple, Dict +import inspect +from pydantic import BaseModel, create_model, ValidationError, parse_obj_as, TypeAdapter +from typing import Any, Optional, get_origin, get_args, Union import fsspec from rich.columns import Columns @@ -150,6 +153,7 @@ class CommandResult: + class BaseCommand(metaclass=SubCommandMeta): """ Base class for command line interfaces. @@ -214,16 +218,86 @@ class BaseCommand(metaclass=SubCommandMeta): """ Handle calls to the OWILocal instance as command invocations. """ return self.do(cmd, *args, **kwargs) + def _cast_args(self, func, args: Tuple[Any, ...], kwargs: Dict[str, Any]) -> Tuple[Tuple[Any, ...], Dict[str, Any]]: + """ + Converts args and kwargs to match the types specified in the function's signature using pydantic. + + Parameters: + func (callable): The function whose signature will be used for type conversion. + args (tuple): The positional arguments to convert. + kwargs (dict): The keyword arguments to convert. + + Returns: + tuple: A tuple containing the converted args and kwargs. + """ + sig = inspect.signature(func) + parameters = list(sig.parameters.values()) + + converted_args = [] + converted_kwargs = {} + + # Convert positional arguments (*args) that match the function signature + for i, arg in enumerate(args): + if i < len(parameters): + param = parameters[i] + expected_type = param.annotation + + if expected_type == inspect.Parameter.empty: + converted_args.append(arg) + else: + try: + # Use TypeAdapter to convert the argument + type_adapter = TypeAdapter(expected_type) + converted_args.append(type_adapter.validate_python(arg)) + except (ValueError, TypeError) as e: + print(f"WARNING - Could not convert arg[{i}]='{arg}' to {expected_type}: {e}") + converted_args.append(arg) + + # Include remaining *args as-is + if len(args) > len(parameters): + converted_args.extend(args[len(parameters):]) + + # Convert keyword arguments (*kwargs) that match the function signature + for name, param in sig.parameters.items(): + if param.kind in (inspect.Parameter.KEYWORD_ONLY, inspect.Parameter.POSITIONAL_OR_KEYWORD): + if name in kwargs: + expected_type = param.annotation + value = kwargs[name] + if expected_type != inspect.Parameter.empty: + try: + type_adapter = TypeAdapter(expected_type) + converted_kwargs[name] = type_adapter.validate_python(value) + except (ValueError, TypeError) as e: + print(f"WARNING - Could not convert kwarg '{name}'='{value}' to {expected_type}: {e}") + converted_kwargs[name] = value + else: + converted_kwargs[name] = value + + # Include any additional kwargs that weren't in the function signature + for k, v in kwargs.items(): + if k not in converted_kwargs: + converted_kwargs[k] = v + + # Ensure no positional argument conflicts with keyword arguments + for i, arg in enumerate(converted_args): + if i < len(parameters): + param_name = parameters[i].name + if param_name in converted_kwargs: + raise TypeError(f"Got multiple values for argument '{param_name}'") + + return tuple(converted_args), converted_kwargs + def do(self, cmd, *args, **kwargs): """ Dispatch to the appropriate sub-command. """ if cmd in self.commands: _param = {"group": str(self.__class__.__name__)} |get_func_params(self.commands[cmd], *args, **kwargs) try: _param["success"] = False + args, converted_kwargs = self._cast_args(self.commands[cmd], args, kwargs) if cmd =="help": - self.console.print(self.commands[cmd](*args, **kwargs)) + self.console.print(self.commands[cmd](*args, **converted_kwargs)) return - returns = self.commands[cmd](*args, **kwargs) + returns = self.commands[cmd](*args, **converted_kwargs) _param["success"] = returns.success if isinstance(returns, CommandResult) else True _param["msg"] = returns.msg if isinstance(returns, CommandResult) else "" except Exception as e: diff --git a/owilix/cmd/query.py b/owilix/cmd/query.py index bc58805..5416778 100644 --- a/owilix/cmd/query.py +++ b/owilix/cmd/query.py @@ -140,7 +140,7 @@ def slice(self, local_specifier, remote_specifier, batch_size=500, prefetch=3, partitioned_by="", access="public", overwrite_ignore=True, import_collection="userslice", chunk_size=1000000, internalID=None, - **kwargs ): + **kwargs): """ Executes the query over the datasets selected by specified local and remote specifiers and applies select and where clause provided in kwargs @@ -246,7 +246,7 @@ def slice(self, local_specifier, remote_specifier, _changelog.append(str(results)) except Exception as e: - self.console.exception(e) + self.console.print_exception() self.console.log(f"Dataset with id {_ds.internalID} at {_ds.path} is most likely corrupted and should be removed.") finally: _ds.append_changelog("Slice Result: "+",".join(_changelog), True) diff --git a/owilix/core/metadata.py b/owilix/core/metadata.py index a68e11a..ae3f2b6 100644 --- a/owilix/core/metadata.py +++ b/owilix/core/metadata.py @@ -574,6 +574,7 @@ class Dataset: _provenance_new = set([f"{create_provenance_url(d, files, select=select, where=where)}" for d in datasets]) _provenance_old = set(self.metadata["provenance"]) _provenance = list(_provenance_old.union(_provenance_new)) + self.metadata["provenance"] = _provenance _overlapping = set([parse_provenance_url(u).get("internalID", None) for u in _provenance_old.intersection(_provenance_new)]) diff --git a/pyproject.toml b/pyproject.toml index e24e42f..a49f0da 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -35,6 +35,7 @@ click= "8.1.6" irods-fsspec="^0.0.1" duckdb = {version = "^1.0.0"} ipython = {version = "^8.24.0"} +pydantic = "^2.5.3" rich = "^13.7.0" numpy = "^1.26.4" -- 2.51.2 From 706b635ff39431f9acd2f12ab9e9067c252883e9 Mon Sep 17 00:00:00 2001 From: Michael Granitzer Date: Mon, 24 Jun 2024 08:05:33 +0200 Subject: [PATCH 3/3] feat: added pyarrow producer / consumer pipeline with opensearch consumer --- docs/source/streams.md | 23 +++ owilix/_version.py | 4 +- owilix/cli.py | 12 +- owilix/cmd/local.py | 2 +- owilix/cmd/query.py | 224 +++++++---------------------- owilix/cmd/remote.py | 28 ++-- owilix/core/duckdb.py | 75 ++-------- owilix/core/metadata.py | 109 +++++++------- owilix/core/repository.py | 5 +- owilix/core/stream.py | 165 +++++++++++++++++++++ owilix/plugins/__init__.py | 0 owilix/plugins/search_consumers.py | 97 +++++++++++++ tests/do_oa.py | 116 +++++++++++++++ tests/test_cli.py | 174 ++++++++++++++++++++++ 14 files changed, 726 insertions(+), 308 deletions(-) create mode 100644 docs/source/streams.md create mode 100644 owilix/core/stream.py create mode 100644 owilix/plugins/__init__.py create mode 100644 owilix/plugins/search_consumers.py create mode 100644 tests/do_oa.py create mode 100644 tests/test_cli.py diff --git a/docs/source/streams.md b/docs/source/streams.md new file mode 100644 index 0000000..938e72e --- /dev/null +++ b/docs/source/streams.md @@ -0,0 +1,23 @@ +# PyArrow Streams + +`owilix` supports creating pyarrow streams from a select / where query and consuming it with different consumers, configured via the command line. +Streaming offers a flexible way to work with OWI data that are either streamed from local files or +remote files. + +Streaming supports the following consumers (which can be again provide the stream to exgternal proceses) + +## Streaming to Stdout + +A standard case is to reuse the stream in another data via stdin. This can be done by using the `owilix.core.stream.ConsumeToStdout` consumer and the following command as example + +```sh +owilix query stream --remote lrz:2023-12-3 select=url,title "where=url_suffix='at'" | nc -l 1234 +``` + +## Streaming to a host:port + +Another common use case is to stream the data to a network port. This can be done by using the `owilix.core.stream.ConsumeToSocket` consumer and the following command as example + +```sh +owilix query stream --local all:2023-12-3 select=url,title "where=url_suffix='at'" consumer="owilix.core.stream.ConsumeToSocket" host="localhost" port=1234 +``` diff --git a/owilix/_version.py b/owilix/_version.py index ca3a9a8..0de2ddb 100644 --- a/owilix/_version.py +++ b/owilix/_version.py @@ -1,3 +1,3 @@ # These version placeholders will be replaced later during substitution. -__version__ = "0.10.0" -__version_tuple__ = (0, 8, 0, "post", 1, "9d9a2b7") +__version__ = "0.10.0-post.2+06b6e1e" +__version_tuple__ = (0, 10, 0, "post", 2, "06b6e1e") diff --git a/owilix/cli.py b/owilix/cli.py index 68c5fc0..b178099 100644 --- a/owilix/cli.py +++ b/owilix/cli.py @@ -198,7 +198,9 @@ def logs(ctx, module): """ if module=="lexis": fn = ctx.obj['OWI'].get_lexis_log_filename() - if not os.path.exists(fn): ctx.obj["CONSOLE"].log(f"Lexis log {fn} not found") + if not os.path.exists(fn): + ctx.obj["CONSOLE"].log(f"Lexis log {fn} not found") + return with open(ctx.obj['OWI'].get_lexis_log_filename(), "r") as log: lines = log.readlines() ctx.obj["CONSOLE"].log(lines) @@ -217,10 +219,14 @@ def logs(ctx, module): ctx.obj["CONSOLE"].print("Unkown module {module}: options are lexis, events, errors") -def main(): - cmds = [local, clean, logs, config, remote, admin, query] +def register_commands(cli): + cmds = [local, clean, logs, config, remote, admin, query] for i in cmds: cli.add_command(i) + + +def main(): + register_commands(cli) cli(obj={}) diff --git a/owilix/cmd/local.py b/owilix/cmd/local.py index 4404025..af39f3f 100644 --- a/owilix/cmd/local.py +++ b/owilix/cmd/local.py @@ -74,7 +74,7 @@ def ls(self, specifier, *args, **kwargs): @LocalCommands.register -def free(self, specifier, *args, **kwargs): +def free(self, specifier): """ removes the datasets locally diff --git a/owilix/cmd/query.py b/owilix/cmd/query.py index 5416778..15fb00e 100644 --- a/owilix/cmd/query.py +++ b/owilix/cmd/query.py @@ -1,3 +1,4 @@ +import importlib import os import sys import uuid @@ -7,7 +8,8 @@ from typing import Dict, List, Optional from fsspec import AbstractFileSystem from owilix.cmd.base import BaseCommand, SubCommand, currentItemProgress, CommandResult, ask_yes_no -from owilix.core.duckdb import OWIlixSQLQuery, OWIDuckDBSelect, OWIDuckDBCopy, OWIDuckDBArrow +from owilix.core.duckdb import OWIlixSQLQuery, OWIDuckDBSelect, OWIDuckDBCopy +from owilix.core.stream import OWIDuckDBArrow, ConsumeToSocket, ConsumeToStdouts from owilix.core.metadata import Dataset, fill_metadata @@ -197,6 +199,7 @@ def slice(self, local_specifier, remote_specifier, _ds = _ds[0] _same = _ds.update_provenance(_datasets, files=files, select=select, where=where) + _ds.metadata.update(kwargs) # udpate metadata with additional kwargs if len(_same)>0 and not ignore_provenance: self.console.print(f"Dataset with id {_ds.internalID} /{_ds.title} already contains the same " f"files and query in its provenance list. Skipping the following datasets " @@ -259,14 +262,14 @@ def slice(self, local_specifier, remote_specifier, @QueryCommands.register def stream(self, local_specifier: str, remote_specifier: str, - select: str = "url,domain_label,title,plain_text", - where: Optional[str] = "", limit: Optional[int] = None, - files: str = "**/*.parquet", verbose: bool = False, - host:str = None, port:int = None, - pq_batch_size: int = 1, batch_size: int = 100, prefetch: int = 2): + select: str = "url,domain_label,title,plain_text", + where: Optional[str] = "", limit: Optional[int] = None, + files: str = "**/*.parquet", verbose: bool = False, + pq_batch_size: int = 1, batch_size: int = 100, prefetch: int = 2, queue_size:int =5, + consumer: str = "owilix.core.stream.ConsumeToStdouts", **kwargs): """ Executes the query over the datasets selected by specified local and remote specifiers - and applies select and where clause provided in kwargs. Provides a stream of arrow data strucures over stdoout + and applies select and where clause provided in kwargs. Provides a stream of arrow data structures over stdout to be consumed in a pipe. Args: @@ -276,178 +279,59 @@ def stream(self, local_specifier: str, remote_specifier: str, where (str, optional): WHERE clause to be applied in the SELECT statement. Defaults to an empty string. limit (int, optional): Limit on the number of rows to return. Defaults to None. files (str): Glob pattern for selecting files in both local and remote locations. Defaults to "**/*.parquet". - explain (bool): Whether to explain the query instead of executing it. Defaults to False. + verbose (bool): Whether to print detailed information to stderr. Defaults to False. pq_batch_size (int): Number of parquet files to consider in one batch. Defaults to 1. batch_size (int): Number of rows per query to be yielded back. Defaults to 100. + queue_size (int): Size of buffer queue in (multiplied by prefetch). Defaults to 5. prefetch (int): Number of batches to prefetch. Defaults to 2. - """ - # todo: print control messages to stderr, but make it configurable - import pyarrow as pa - import pyarrow.ipc as ipc - import queue - import threading - import sys - - def producer(generator, buffer_queue): - for batch in generator: - buffer_queue.put(batch) # Put the batch in the queue - buffer_queue.put(None) # Signal the end of the stream - - def stream_to_socket(buffer_queue, host, port): - with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as sock: - sock.connect((host, port)) - print(f"Connected to {host}:{port}") - - schema_written = False - writer = None - while True: - batch = buffer_queue.get() - if batch is None: - break # End of stream - - if not schema_written: - # Send schema as the first message - writer = ipc.RecordBatchStreamWriter(sock.makefile('wb'), batch.schema) - schema_written = True - - # Write the batch to the socket - writer.write_table(batch) - buffer_queue.task_done() - - # Close the writer after processing all batches - if writer: - writer.close() - print("Finished sending data.") - - def stream_to_output(buffer_queue, output): - schema_written = False - writer = None # Initialize the writer as None - while True: - batch = buffer_queue.get() - if batch is None: - break # End of stream - - if not schema_written: - # Initialize the writer with the schema from the first batch - writer = ipc.RecordBatchStreamWriter(output, batch.schema) - schema_written = True - - writer.write_table(batch) - buffer_queue.task_done() - - # Close the writer after processing all batches - if writer: - writer.close() - - def configure_kafka_producer(): - from confluent_kafka import Producer - import json - - # Kafka Configuration - KAFKA_BOOTSTRAP_SERVERS = 'localhost:9092' - KAFKA_TOPIC = 'pyarrow_stream' - - # Kafka Producer configuration - producer_config = { - 'bootstrap.servers': KAFKA_BOOTSTRAP_SERVERS, - 'client.id': 'pyarrow-stream-producer' - } - - # Create a Kafka producer - return Producer(producer_config) - - def stream_to_kafka(buffer_queue, topic): - while True: - batch = buffer_queue.get() - if batch is None: - break # End of stream - - # Serialize the batch to a bytes buffer - sink = pa.BufferOutputStream() - writer = ipc.RecordBatchStreamWriter(sink, batch.schema) - writer.write_table(batch) - writer.close() - - # Get the serialized data as bytes - data = sink.getvalue().to_pybytes() - - # Send serialized data to Kafka - producer.produce(topic, data) - producer.flush() # Ensure the message is sent - - buffer_queue.task_done() + consumer (str): The fully qualified class name of the consumer to use for processing batches. Defaults to "owilix.core.stream.ConsumeToStdouts". + **kwargs: Additional parameters to be passed to the consumer's constructor. + Raises: + ImportError: If the consumer class cannot be imported. + AttributeError: If the consumer class does not exist in the specified module. + + Examples: + - local stream to stdout: + owilix query stream --local all:2023-12-3 select=url,title "where=url_suffix='at'" + - local stream to be written to a network port: + owilix query stream --local all:2023-12-3 select=url,title "where=url_suffix='at'" consumer="owilix.core.stream.ConsumeToSocket" host="localhost" port=1234 + """ + # Load the consumer class dynamically + module_name, class_name = consumer.rsplit('.', 1) + try: + module = importlib.import_module(module_name) + ConsumerClass = getattr(module, class_name) + except ImportError as e: + print(f"Error: Failed to import module '{module_name}'.", file=sys.stderr) + raise e + except AttributeError as e: + print(f"Error: Module '{module_name}' does not have a class '{class_name}'.", file=sys.stderr) + raise e + + # Instantiate the consumer class with any additional keyword arguments + _, kwargs = self._cast_args(ConsumerClass.__init__, [], kwargs) + consumer_instance = ConsumerClass(**kwargs) + + # Retrieve the list of files and datasets to process all_files = self.get_all_files(files, local_specifier, remote_specifier, print_it=False) _datasets = all_files.keys() all_files = self.group_all_files_by_fs(all_files) + if verbose: print(f"Found '{sum([len(v) for v in all_files.values()])}' parquet files " - f"in {len(all_files.keys())} filesystems over {len(_datasets)}. " - f"Running queries against them and providing results as binary backpressure controlled output stream.", file=sys.stderr) + f"in {len(all_files.keys())} filesystems over {len(_datasets)} datasets. " + f"Running queries against them and providing results as binary backpressure-controlled output stream.", file=sys.stderr) + + # Prepare the SQL query sql = (OWIlixSQLQuery.from_templates("pq_select") .select(select) .where(where) - .limit(limit)) - - db = OWIDuckDBArrow(all_files, sql, explain=False, - pq_batch_size=pq_batch_size, - batch_size=batch_size, - prefetch=prefetch) - buffer_queue = queue.Queue(maxsize=5*prefetch) - - # Start the producer thread - producer_thread = threading.Thread(target=producer, args=(db.query_aggregator(), buffer_queue)) - producer_thread.start() - - if host and port: - import socket - with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as sock: - sock.connect((host, port)) - stream_to_output(buffer_queue, sock.makefile('wb')) - else: - # Stream the result to stdout (Linux pipe) - with open(sys.stdout.fileno(), 'wb', closefd=False) as output: - stream_to_output(buffer_queue, output) - - # Wait for the producer to finish - producer_thread.join() - buffer_queue.join() - -@QueryCommands.register -def consume(self): - """ - Consumes an arrow stream from stdin and prints the record batches to the console. Used for testing pyarrow stream command - :param self: - :return: - """ - import pyarrow as pa - import pyarrow.ipc as ipc - import sys - from rich.console import Console - from rich.table import Table - - def print_batch(batch): - # Create a rich Table to display the batch - table = Table(title="Arrow Stream Batch") - - # Add columns to the table - for column in batch.schema.names: - table.add_column(column) - - # Add rows to the table - for row in batch.to_pylist(): - table.add_row(*(str(row[col]) for col in batch.schema.names)) - - # Print the table to the console - self.console.print(table) - - # Open stdin as a file-like object for reading binary data - input_stream = sys.stdin.buffer - - # Create a RecordBatchStreamReader to read the PyArrow stream - reader = ipc.RecordBatchStreamReader(input_stream) - - # Process each record batch in the stream - for i in range(reader.num_record_batches): - batch = reader.get_batch(i) - print_batch(batch) \ No newline at end of file + .limit(str(limit))) + + # Initialize and start the streaming process with the consumer instance + db_instance = OWIDuckDBArrow(all_files, sql, explain=False, + pq_batch_size=pq_batch_size, + batch_size=batch_size, + prefetch=prefetch) + db_instance.stream(consumer_instance, queue_size) diff --git a/owilix/cmd/remote.py b/owilix/cmd/remote.py index b50a4be..e893330 100644 --- a/owilix/cmd/remote.py +++ b/owilix/cmd/remote.py @@ -85,22 +85,20 @@ def ls(self, specifier, *args, details=False, **kwargs): return CommandResult(success=True, object=datasets, msg=f"Showed {len(datasets)} datasets") @RemoteCommands.register -def pull(self, specifier, *args, **kwargs): +def pull(self, specifier, files = "**/*", overwrite = False): """"pulls the datasets with the given specifier to the local repository Args: overwrite(bool) : if True, files are overwritten if they exists (default False) files(str): a glob pattern to select files within the datasets to be pulled """ - file_select = kwargs.get("files", "**/*") - overwrite = kwargs.get("overwrite", False) self.console.print(f"Fetching datasets for specifier {specifier}") datasets = self.list_remote_dataset_by_specifier(specifier) self.show_datasets(datasets) if self.autoyes or ask_yes_no(self.console, f"Download these datasets:"): for d in datasets: self.console.print(f"Fetching files for {d.title}") - _remote_files = [rel_can_path(_p, d.path) for _p in self.owi.remote_data.files(d, file_select)] + _remote_files = [rel_can_path(_p, d.path) for _p in self.owi.remote_data.files(d, files)] _ds_local = self.owi.local.list(access=d.access, query={"internalID": d.internalID, "collectionName": d.collectionName}) if len(_ds_local)==0: @@ -109,7 +107,7 @@ def pull(self, specifier, *args, **kwargs): raise ValueError(f"Multiple datasets found for {d.internalID}. Dataset inconsistent") self.console.print(f"Found {len(_remote_files)} remote files. Syncing with local files") _local_files = [rel_can_path(_p, _ds_local[0].path) - for _p in self.owi.local.files(_ds_local[0], file_select)] + for _p in self.owi.local.files(_ds_local[0], files)] _missing_files = set(_remote_files) - set(_local_files) if not overwrite else set(_remote_files) if len(_missing_files)==0: self.console.print(f"Dataset {d.title} is already up to date. Missing files are {len(_missing_files)}.") @@ -138,7 +136,7 @@ def doctor(self, *args, **kwargs): @RemoteCommands.register -def push(self, specifier, *args, **kwargs): +def push(self, specifier, files = "**/*", mdupdate = True, overwrite = False, **kwargs): """" push local datasets fitting the specifier to the local repository. Pushing may change the internalID of the local dataset if the dataset does not exist on the server. @@ -154,9 +152,6 @@ def push(self, specifier, *args, **kwargs): dataCenter(str): specify the dataCenter the dataset is created in. HNote that only works if the dataset is not already in a different, known datacenter. overwrites the dataCenter specified in the selector. """ - file_select = kwargs.get("files", "**/*") - md_update = kwargs.get("mdupdate", True) - overwrite = kwargs.get("overwrite", False) dataCenter= kwargs.get("dataCenter", self.owi.parse_specifier(specifier).get("data_center",None)) self.console.print(f"Fetching datasets for specifier {specifier}") @@ -165,7 +160,7 @@ def push(self, specifier, *args, **kwargs): if self.autoyes or ask_yes_no(self.console, f"Upload these datasets (overwrite:{overwrite}):"): for d in datasets: self.console.print(f"Fetching files for {d.title}") - _local_files = [rel_can_path(_p, d.path) for _p in self.owi.local.files(d, file_select)] + _local_files = [rel_can_path(_p, d.path) for _p in self.owi.local.files(d, files)] _ds_remote = self.owi.remote_data.list(access=d.access, query={"internalID": d.internalID, "collectionName": d.collectionName}) if len(_ds_remote)==0: @@ -184,7 +179,7 @@ def push(self, specifier, *args, **kwargs): elif len(_ds_remote)>1: raise ValueError(f"Multiple datasets found for {d.internalID}. Dataset inconsistent") else: - if md_update: + if mdupdate: _ds_remote[0].metadata.update(cast_metadata({k:v for k,v in d.metadata.items() if k!="dataCenter"})) _ds_remote[0].reformat_metadata() self.owi.remote_data.update_metadata(_ds_remote[0]) @@ -195,7 +190,7 @@ def push(self, specifier, *args, **kwargs): self.console.print(f"Found {len(_local_files)} local files. " f"Syncing with remote files at {_ds_remote[0].dataCenter}") _remote_files = [rel_can_path(_p, _ds_remote[0].path) - for _p in self.owi.remote_data.files(_ds_remote[0], file_select)] + for _p in self.owi.remote_data.files(_ds_remote[0], files)] _missing_files = set(_local_files) - set(_remote_files) if not overwrite else set(_local_files) if len(_missing_files)==0: self.console.print(f"Dataset {d.title} is already up to date") @@ -219,14 +214,13 @@ def push(self, specifier, *args, **kwargs): return CommandResult(success=True, object=datasets, msg=f"Pushed {len(datasets)} datasets") @RemoteCommands.register -def diff(self, specifier, *args, **kwargs): +def diff(self, specifier, files=None, **kwargs): """ runs a dataset and/or file-level diff Args: files(str): a glob pattern to select files for. Use **/* for all files. If not set, diff will only be applied on the dataset level """ - file_select = kwargs.get("files", None) self.console.print(f"Fetching local and remote datasets for specifier {specifier}") _local_datasets = self.list_local_datasets_by_specifier(specifier) _remote_datasets = self.list_remote_dataset_by_specifier(specifier) @@ -241,7 +235,7 @@ def diff(self, specifier, *args, **kwargs): self.console.print("Dataets available only remotely:") self.show_datasets([_d for _d in _remote_datasets if _d.internalID not in _ids]) - if file_select is None: + if files is None: self.console.print(f"Diff done (for diff on a file level specify files=**/*") return @@ -257,8 +251,8 @@ def diff(self, specifier, *args, **kwargs): for i in _diff_print: i["key"] = i["key"] if i["local"]==i["remote"] else "[warning]"+str(i["key"])+"[/warning]" self.show_table(_diff_print, order="key,local,remote") # now check for file diffs - _local_files = [rel_can_path(_p, _localds.path) for _p in self.owi.local.files(_localds, file_select)] - _remote_files = [rel_can_path(_p, _remoteds.path) for _p in self.owi.remote_data.files(_remoteds, file_select)] + _local_files = [rel_can_path(_p, _localds.path) for _p in self.owi.local.files(_localds, files)] + _remote_files = [rel_can_path(_p, _remoteds.path) for _p in self.owi.remote_data.files(_remoteds, files)] self.console.print(f"Found {len(_local_files)} local files, {len(_remote_files)} remote files. ") _lonly = set(_local_files) - set(_remote_files) if len(_lonly)>1: diff --git a/owilix/core/duckdb.py b/owilix/core/duckdb.py index ce50817..00967e6 100644 --- a/owilix/core/duckdb.py +++ b/owilix/core/duckdb.py @@ -350,6 +350,7 @@ class OWIDuckDBSelect: self.as_dict = as_dict self.max_mem = max_mem self.retry_count = retry_count + self.logger = logging.getLogger("owilix") def retry_then_raise(self, connenction, create_owi_slice_query): @@ -396,11 +397,11 @@ class OWIDuckDBSelect: cursor = self.retry_then_raise(conn, _query.sql) columns = None - _logger.debug(f"Connection to filesystem '{fs.protocol}' for {len(pq_batch.files)} files opened.") + self.logger.debug(f"Connection to filesystem '{fs.protocol}' for {len(pq_batch.files)} files opened.") while True: - _logger.debug(f"Fetching {batch_size} rows from {fs.protocol} for {len(pq_batch.files)} files") + self.logger.debug(f"Fetching {batch_size} rows from {fs.protocol} for {len(pq_batch.files)} files") results = cursor.fetchmany(batch_size) - _logger.debug(f"Fetch of {batch_size} rows done from {fs.protocol} for {len(pq_batch.files)} files") + self.logger.debug(f"Fetch of {batch_size} rows done from {fs.protocol} for {len(pq_batch.files)} files") if not results: break if self.as_dict: @@ -410,12 +411,12 @@ class OWIDuckDBSelect: yield results except Exception as e: - _logger.exception(f"Error when executing {query.sql} on files {pq_batch}") + self.logger.exception(f"Error when executing {query.sql} on files {pq_batch}") raise e finally: if conn: conn.close() - _logger.debug(f"Connection to {fs.protocol} for {len(pq_batch.files)} files closed.") + self.logger.debug(f"Connection to {fs.protocol} for {len(pq_batch.files)} files closed.") tmp_dir.cleanup() def query_aggregator(self) -> Generator[List[tuple], None, None]: @@ -467,7 +468,7 @@ class OWIDuckDBSelect: for result_batch in future.result(): yield result_batch except Exception as e: - _logger.exception(f"Error processing task {task}: {e}") + self.logger.exception(f"Error processing task {task}: {e}") class OWIDuckDBCopy (OWIDuckDBSelect): @@ -543,7 +544,7 @@ class OWIDuckDBCopy (OWIDuckDBSelect): DROP TABLE IF EXISTS owi_slice; CREATE TABLE owi_slice AS """ + query.files([f[0] for f in pq_batch.files]).sql # could be also done in smaller batches - _logger.debug(f"Connection to filesystem '{fs.protocol}' for {len(pq_batch.files)} files opened.") + self.logger.debug(f"Connection to filesystem '{fs.protocol}' for {len(pq_batch.files)} files opened.") # Execute the query to create owi_slice self.retry_then_raise(conn, create_owi_slice_query) @@ -597,7 +598,7 @@ class OWIDuckDBCopy (OWIDuckDBSelect): conn.close() except Exception as e: - _logger.exception(f"Error when executing {query.sql} on files {pq_batch}") + self.logger.exception(f"Error when executing {query.sql} on files {pq_batch}") _error=str(e) finally: @@ -605,62 +606,6 @@ class OWIDuckDBCopy (OWIDuckDBSelect): yield [{"message": _error, "group":_sub_path, "num_files":len(pq_batch.files), "count":0, "success":0}] if conn: conn.close() - _logger.debug(f"Connection to {fs.protocol} for {len(pq_batch.files)} files closed.") + self.logger.debug(f"Connection to {fs.protocol} for {len(pq_batch.files)} files closed.") temp_dir.cleanup() -class OWIDuckDBArrow(OWIDuckDBSelect): - """ - A class to execute a DuckDB select, similar to OWIDuckDBSelect, but using Apache Arrow as return results - - Attributes: - see OWIDuckDBSelect - """ - - def run_query(self, fs: AbstractFileSystem, - pq_batch: ParquetBatch, - query: OWIlixSQLQuery, - batch_size: int = 0) -> Generator[List[tuple], None, None]: - """ - Run a SQL query on a batch of parquet files using DuckDB. - - Args: - fs (AbstractFileSystem): The filesystem containing the parquet files. - pq_batch (List[str]): A batch of parquet file paths to query. - query (OWIlixSQLQuery): The SQL query to execute. - batch_size (int, optional): The number of rows to fetch in each batch. Defaults to self.batch_size. - - Yields: - Generator[List[tuple], None, None]: A generator yielding batches of query results. - """ - if batch_size <= 0: - batch_size = self.batch_size - conn, tmp_dir = None, tempfile.TemporaryDirectory() - try: - temp_db_path = os.path.join(tmp_dir.name, 'temp_owi_select_duckdb.db') - conn = duckdb.connect(database=temp_db_path) - conn.execute(f"PRAGMA memory_limit='{self.max_mem}'") - conn.register_filesystem(fs) - _query = query.files([f[0] for f in pq_batch.files]) if len(pq_batch.files) > 0 else query - # Format the SQL query using the query_args from the ParquetBatch - if pq_batch.query_args: - _query = _query.format(**pq_batch.query_args) - - cursor = self.retry_then_raise(conn,_query.sql) - _logger.debug(f"Connection to filesystem '{fs.protocol}' for {len(pq_batch.files)} files opened.") - while True: - _logger.debug(f"Fetching {batch_size} rows from {fs.protocol} for {len(pq_batch.files)} files") - results = cursor.fetch_arrow_table(batch_size) - _logger.debug(f"Fetch of {batch_size} rows done from {fs.protocol} for {len(pq_batch.files)} files") - if not results or results.num_rows == 0: - break - yield results - - except Exception as e: - _logger.exception(f"Error when executing {query.sql} on files {pq_batch}") - raise e - finally: - if conn: - conn.close() - _logger.debug(f"Connection to {fs.protocol} for {len(pq_batch.files)} files closed.") - tmp_dir.cleanup() - diff --git a/owilix/core/metadata.py b/owilix/core/metadata.py index ae3f2b6..7ea46e5 100644 --- a/owilix/core/metadata.py +++ b/owilix/core/metadata.py @@ -1,8 +1,9 @@ -import datetime + import json import re import uuid from collections import defaultdict +from dateutil import parser from dataclasses import dataclass from datetime import datetime, timedelta import os @@ -417,24 +418,34 @@ class Dataset: """ Represents a dataset in the repository. Metadata are set in the .metadata property and can be accessed as attributes. - Workflow relevant properties are available as atttirbutes (e.g. path) - todo: integrate metadta checks and validation in dataset. + Workflow relevant properties are available as attributes (e.g. path). + TODO: integrate metadata checks and validation in dataset. """ def __init__(self, repository, path, **metadata): self.repository = repository self.path = path - self.metadata = cast_metadata(metadata) + self._metadata = cast_metadata(metadata) # Use a private attribute to store metadata self._change_log = None - if "internalID" not in metadata: + + if "internalID" not in self._metadata: raise ValueError("Dataset metadata must contain an internalID") - if "title" not in metadata: - self.metadata["title"] = "UNKNOWN TITLE" + if "title" not in self._metadata: + self._metadata["title"] = "UNKNOWN TITLE" + # Property to encapsulate metadata + @property + def metadata(self): + return self._metadata + + @metadata.setter + def metadata(self, value): + self._metadata = cast_metadata(value) + self.update_lastchanged() def __getattr__(self, name): - if name in self.metadata: - return self.metadata[name] + if name in self._metadata: + return self._metadata[name] else: raise AttributeError(f"Dataset has no attribute {name}") @@ -444,19 +455,19 @@ class Dataset: @property def startDate(self): try: - return datetime.strptime(self.metadata.get("startDate", None), '%Y-%m-%d') + return parser.parse(self._metadata.get("startDate", None)) except: return None @property def endDate(self): try: - return datetime.strptime(self.metadata.get("endDate", None), '%Y-%m-%d') + return parser.parse(self._metadata.get("endDate", None)) except: return None def get_changelog(self): - if self._change_log==None: + if self._change_log is None: _content = self.repository.readlines(self, "changelog.json") self._change_log = json.loads(_content) if _content else None self._change_log = self._change_log if self._change_log is not None else [] @@ -469,15 +480,16 @@ class Dataset: "msg": msg }) self._change_log = _change_log - if save: self.save_changelog() + if save: + self.save_changelog() def save_changelog(self): - if self._change_log!=None: + if self._change_log is not None: self.repository.writelines(self, "changelog.json", json.dumps(self._change_log)) return self._change_log def update_lastchanged(self): - self.metadata['lastChanged'] = self._get_now_formatted() + self._metadata['lastChanged'] = self._get_now_formatted() def _get_now_formatted(self): return datetime.now().strftime("%Y-%m-%d %H:%M:%S.%f") @@ -487,27 +499,28 @@ class Dataset: Reformat the metadata to a dictionary """ self.metadata["title"] = reformat_metadata(self.metadata)["title"] + self.update_lastchanged() def consolidate_metadata(self, infer=False, **kwargs): """ Consolidate metadata by checking and updating the metadata fields. If infer is True, metadata is inferred from the files in the dataset and updated. - kwargs contains a set of metadata valus to be updated. Key is the metadata field name. - value is either the value to be set, None for deleting the metadata entry or (default) for getting the default value + kwargs contains a set of metadata values to be updated. Key is the metadata field name. + Value is either the value to be set, None for deleting the metadata entry, or (default) for getting the default value. """ - _defaults= fill_metadata({}, overwrite_inferred=False) - if infer: self.metadata.update(infer_metadata_from_files(self.repository.files_details(self))) - for k,v in _defaults.items(): + _defaults = fill_metadata({}, overwrite_inferred=False) + if infer: + self.metadata.update(infer_metadata_from_files(self.repository.files_details(self))) + for k, v in _defaults.items(): if k not in self.metadata: self.metadata[k] = v for key, value in kwargs.items(): if not value: - if key in self.metadata: del self.metadata[key] + if key in self.metadata: + del self.metadata[key] else: self.metadata[key] = value if value != "default" else _defaults[key] - self.metadata = cast_metadata(self.metadata) - - + self.metadata = cast_metadata(self.metadata) # This will trigger the setter and update_lastchanged def filter_by_date_range(self, start_date=None, end_date=None): start = self.startDate @@ -522,29 +535,30 @@ class Dataset: return True # No filtering @staticmethod - def filter_datasets(_datasets, day, duration, query)-> list: - - # query filter + def filter_datasets(_datasets, day, duration, query) -> list: + # Query filter if query is not None and query != {}: def matches_query(_d): _all = [] for key, value in query.items(): regex = key.endswith("*") key = key if not regex else key[:-1] - if not hasattr(_d, key): _all.append(False) + if not hasattr(_d, key): + _all.append(False) if not regex: _all.append(getattr(_d, key) == value) else: - _all.append(re.match(value, getattr(_d,key))) + _all.append(re.match(value, getattr(_d, key))) return all(_all) _datasets = [_d for _d in _datasets if matches_query(_d)] - if len(_datasets)==0: return [] - # time filter + if len(_datasets) == 0: + return [] + # Time filter if day == "latest": - day = max( _d.startDate for _d in _datasets) + day = max(_d.startDate for _d in _datasets) elif day is not None and not isinstance(day, datetime): day = datetime.strptime(day, '%Y-%m-%d') @@ -554,24 +568,23 @@ class Dataset: return [_d for _d in _datasets if _d.filter_by_date_range(startdate, day)] - @staticmethod def from_lexis_http(repository, path, **metadata): def convert_value(k, v): if isinstance(v, list) and len(v) == 1: return v[0] return v + lower_case = ["AlternateIdentifier", "CreationDate", "CustomMetadataSchema", "RelatedSoftware"] for l in lower_case: if l in metadata: - metadata[l[0].lower()+l[1]] = metadata.pop(l) - _md = extract_metadata_from_title(metadata["metadata"].get("title",["UKNOWN TITLE"])[0]) + metadata[l[0].lower() + l[1]] = metadata.pop(l) + _md = extract_metadata_from_title(metadata["metadata"].get("title", ["UKNOWN TITLE"])[0]) _md = _md | metadata["flags"] | metadata["location"] | {k: convert_value(k, v) for k, v in metadata["metadata"].items()} return Dataset(repository, path, **_md) - def update_provenance(self, datasets: List, files:str = None, select:str = None, - where:str=None): - _provenance_new = set([f"{create_provenance_url(d, files, select=select, where=where)}" for d in datasets]) + def update_provenance(self, datasets: List, files: str = None, select: str = None, where: str = None): + _provenance_new = set([f"{create_provenance_url(d, files, select=select, where=where)}" for d in datasets]) _provenance_old = set(self.metadata["provenance"]) _provenance = list(_provenance_old.union(_provenance_new)) self.metadata["provenance"] = _provenance @@ -581,19 +594,19 @@ class Dataset: return [d for d in datasets if d.internalID in _overlapping] @staticmethod - def merge_into_new(repository, datasets: List, files:str = None, select:str = None, - where:str=None, access:ACCESSTYPES="project", - collectionName:str="userslice", **kwargs): + def merge_into_new(repository, datasets: List, files: str = None, select: str = None, + where: str = None, access: str = "project", collectionName: str = "userslice", **kwargs): """ - create a new dataset by merging the provided ones and applying files, select and where filters (for updating the provenance) + Create a new dataset by merging the provided ones and applying files, select and where filters (for updating the provenance) """ _new_md = fill_metadata({}, None, False) - _new_md["provenance"] = [ f"{create_provenance_url(d, files, select=select, where=where)}" for d in datasets] + _new_md["provenance"] = [f"{create_provenance_url(d, files, select=select, where=where)}" for d in datasets] _new_md.update(kwargs) - _new_md["access"]=access - _new_md["collectionName"]=collectionName - if "description" not in kwargs: _new_md["description"] = (f"Merged dataset from {len(datasets)} datasets at " - f"{datetime.now()} using parameters files={files}, " - f"select={select}, where={where}") + _new_md["access"] = access + _new_md["collectionName"] = collectionName + if "description" not in kwargs: + _new_md["description"] = (f"Merged dataset from {len(datasets)} datasets at " + f"{datetime.now()} using parameters files={files}, " + f"select={select}, where={where}") return repository.create(**_new_md) diff --git a/owilix/core/repository.py b/owilix/core/repository.py index c9da7b5..2e5b3c8 100644 --- a/owilix/core/repository.py +++ b/owilix/core/repository.py @@ -609,7 +609,8 @@ class FileBasedRepository(AbstractRepository): if len(_md_errors)>0: logger.warning(f"Metadata inconsistencies for dataset {dataset.internalID}: "+",".join(_md_errors.values())) try: - dataset.update_lastchanged() + if not self.fs.exists(os.path.dirname(_p)): + self.fs.mkdirs(os.path.dirname(_p)) with self.fs.open(_p+".json", "w") as _fh: json.dump(dataset.metadata, _fh) except Exception as e: @@ -741,7 +742,7 @@ class IRODSRepository(FileBasedRepository): _md_errors = validate_metadata(dataset.metadata) if len(_md_errors)>0: logger.warning(f"Metadata inconsistencies for dataset {dataset.internalID}: "+",".join(_md_errors.values())) - dataset.update_lastchanged() + return update_metadata_for_irods_collection(self.session.collections.get(_p), dataset.metadata) def change_id(self, dataset, new_id): diff --git a/owilix/core/stream.py b/owilix/core/stream.py new file mode 100644 index 0000000..afd85ef --- /dev/null +++ b/owilix/core/stream.py @@ -0,0 +1,165 @@ +import os +import queue +import socket +import tempfile +import threading +from typing import Generator, List, Callable + +import duckdb +import pyarrow.ipc as ipc +from abc import ABC, abstractmethod +import pyarrow as pa +from fsspec import AbstractFileSystem + +from owilix.core.duckdb import OWIDuckDBSelect, ParquetBatch, OWIlixSQLQuery +import sys +import pyarrow.ipc as ipc + + + +class Consumer(ABC): + """An abstract class for consuming data from a stream of pyarrow table batches.""" + @abstractmethod + def consume(self, batch: pa.Table): + """Process a batch of data.""" + pass + + @abstractmethod + def close(self): + """Perform any cleanup necessary.""" + pass + +class ConsumeToSocket(Consumer): + """ + Class to consume data from a stream of pyarrow table batches and send it over a socket. + """ + def __init__(self, host: str, port: int): + self.host = host + self.port = port + self.socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + self.socket.connect((self.host, self.port)) + self.output = self.socket.makefile('wb') + self.writer = None + + def consume(self, batch: pa.Table): + if self.writer is None and isinstance(batch, pa.Table): + self.writer = ipc.RecordBatchStreamWriter(self.output, batch.schema) + + if self.writer and isinstance(batch, pa.Table): + self.writer.write_table(batch) + else: + print("Warning: The provided batch is not a valid pyarrow Table.", file=sys.stderr) + + def close(self): + if self.writer: + self.writer.close() + self.output.close() + self.socket.close() + + +class ConsumeToStdouts(Consumer): + def __init__(self, **kwargs): + self.output = open(sys.stdout.fileno(), 'wb', closefd=False) + self.writer = None + + def consume(self, batch: pa.Table): + if self.writer is None and isinstance(batch, pa.Table): + self.writer = ipc.RecordBatchStreamWriter(self.output, batch.schema) + + if self.writer and isinstance(batch, pa.Table): + self.writer.write_table(batch) + else: + print("Warning: The provided batch is not a valid pyarrow Table.", file=sys.stderr) + + def close(self): + if self.writer: + self.writer.close() + self.output.close() + + +class OWIDuckDBArrow(OWIDuckDBSelect): + """ + A class to execute a DuckDB select, similar to OWIDuckDBSelect, but using Apache Arrow as return results + + Attributes: + see OWIDuckDBSelect + """ + + def run_query(self, fs: AbstractFileSystem, + pq_batch: ParquetBatch, + query: OWIlixSQLQuery, + batch_size: int = 0) -> Generator[List[tuple], None, None]: + """ + Run a SQL query on a batch of parquet files using DuckDB. + + Args: + fs (AbstractFileSystem): The filesystem containing the parquet files. + pq_batch (List[str]): A batch of parquet file paths to query. + query (OWIlixSQLQuery): The SQL query to execute. + batch_size (int, optional): The number of rows to fetch in each batch. Defaults to self.batch_size. + + Yields: + Generator[List[tuple], None, None]: A generator yielding batches of query results. + """ + if batch_size <= 0: + batch_size = self.batch_size + conn, tmp_dir = None, tempfile.TemporaryDirectory() + try: + temp_db_path = os.path.join(tmp_dir.name, 'temp_owi_select_duckdb.db') + conn = duckdb.connect(database=temp_db_path) + conn.execute(f"PRAGMA memory_limit='{self.max_mem}'") + conn.register_filesystem(fs) + _query = query.files([f[0] for f in pq_batch.files]) if len(pq_batch.files) > 0 else query + # Format the SQL query using the query_args from the ParquetBatch + if pq_batch.query_args: + _query = _query.format(**pq_batch.query_args) + + cursor = self.retry_then_raise(conn,_query.sql) + self.logger.debug(f"Connection to filesystem '{fs.protocol}' for {len(pq_batch.files)} files opened.") + while True: + self.logger.debug(f"Fetching {batch_size} rows from {fs.protocol} for {len(pq_batch.files)} files") + results = cursor.fetch_arrow_table(batch_size) + self.logger.debug(f"Fetch of {batch_size} rows done from {fs.protocol} for {len(pq_batch.files)} files") + if not results or results.num_rows == 0: + break + yield results + + except Exception as e: + self.logger.exception(f"Error when executing {query.sql} on files {pq_batch}") + raise e + finally: + if conn: + conn.close() + self.logger.debug(f"Connection to {fs.protocol} for {len(pq_batch.files)} files closed.") + tmp_dir.cleanup() + + def producer(self, buffer_queue): + for batch in self.query_aggregator(): + buffer_queue.put(batch) # Put the batch in the queue + buffer_queue.put(None) # Signal the end of the stream + + def stream(self, consumer: Consumer, queue_size=5): + """ + starts streaming the query, buffer them via a queue of size queue_size*self.prefetch and consume them using the consumer + """ + buffer_queue = queue.Queue(maxsize=queue_size*self.prefetch) + # Start the producer thread + producer_thread = threading.Thread(target=self.producer, args=(buffer_queue,)) + producer_thread.start() + try: + # Consumer logic + while True: + batch = buffer_queue.get() + if batch is None: + break # End of stream + + consumer.consume(batch) + buffer_queue.task_done() + + # Wait for the producer to finish + producer_thread.join() + buffer_queue.join() + + finally: + # Perform any cleanup needed by the consumer + consumer.close() \ No newline at end of file diff --git a/owilix/plugins/__init__.py b/owilix/plugins/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/owilix/plugins/search_consumers.py b/owilix/plugins/search_consumers.py new file mode 100644 index 0000000..bd0a092 --- /dev/null +++ b/owilix/plugins/search_consumers.py @@ -0,0 +1,97 @@ +import os +import json +import logging +from typing import Any +from opensearchpy import OpenSearch, helpers +import pyarrow as pa + +from owilix.core.stream import Consumer + + +class OpenSearchConsumer(Consumer): + def __init__(self, index_name: str, **kwargs): + """ + Initializes the OpenSearch consumer. + + Args: + index_name (str): The name of the index where data will be inserted. + **kwargs: Additional configuration parameters for OpenSearch. + These can include 'host', 'port', 'timeout', etc. + """ + self.index_name = index_name + self.host = kwargs.get('host', 'localhost') + self.port = kwargs.get('port', 9200) + self.scheme = kwargs.get('scheme', 'http') + self.timeout = kwargs.get('timeout', 30) + + self.username = os.getenv('OWIXILX_OPENSEARCH_USERNAME', kwargs.get('username', None)) + self.password = os.getenv('OWIXLIS_OPENSEARCH_PASSWORD', kwargs.get('password', None)) + + self.client = self._create_opensearch_client() + + self.log_file = kwargs.get('log_file', 'opensearch_errors.log') + logging.basicConfig(filename=self.log_file, level=logging.ERROR) + + self._ensure_index_exists() + + def _create_opensearch_client(self) -> OpenSearch: + """ + Creates an OpenSearch client using provided credentials and configurations. + """ + auth = None + if self.username and self.password: + auth = (self.username, self.password) + + return OpenSearch( + hosts=[{'host': self.host, 'port': self.port}], + http_auth=auth, + use_ssl=(self.scheme == 'https'), + timeout=self.timeout + ) + + def _ensure_index_exists(self): + """ + Checks if the specified index exists in OpenSearch. Creates the index if it does not exist. + """ + if not self.client.indices.exists(index=self.index_name): + self.client.indices.create(index=self.index_name) + print(f"Index '{self.index_name}' created in OpenSearch.") + + def consume(self, batch: pa.Table): + """ + Consumes a batch of data and pushes it to OpenSearch. + + Args: + batch (pa.Table): The batch of data to be inserted. + """ + records = batch.to_pylist() + actions = [ + { + "_index": self.index_name, + "_source": record + } + for record in records + ] + + try: + helpers.bulk(self.client, actions) + except Exception as e: + logging.error(f"Failed to insert records into OpenSearch: {e}") + self._log_failed_records(actions) + + def _log_failed_records(self, actions: Any): + """ + Logs the records that failed to be inserted into OpenSearch to a log file. + + Args: + actions (Any): The list of records that failed to insert. + """ + with open(self.log_file, 'a') as log_file: + for action in actions: + log_file.write(json.dumps(action) + '\n') + + def close(self): + """ + Performs any necessary cleanup. For OpenSearch, no explicit cleanup is required. + """ + pass diff --git a/tests/do_oa.py b/tests/do_oa.py new file mode 100644 index 0000000..1697cef --- /dev/null +++ b/tests/do_oa.py @@ -0,0 +1,116 @@ +import inspect +from pydantic import BaseModel, create_model, ValidationError, TypeAdapter +from typing import Optional, Any, Tuple, Dict + + +def convert_args_kwargs_with_pydantic(func, args: Tuple[Any, ...], kwargs: Dict[str, Any]) -> Tuple[ + Tuple[Any, ...], Dict[str, Any]]: + """ + Converts args and kwargs to match the types specified in the function's signature using pydantic. + + Parameters: + func (callable): The function whose signature will be used for type conversion. + args (tuple): The positional arguments to convert. + kwargs (dict): The keyword arguments to convert. + + Returns: + tuple: A tuple containing the converted args and kwargs. + """ + sig = inspect.signature(func) + parameters = list(sig.parameters.values()) + + converted_args = [] + converted_kwargs = {} + + # Convert positional arguments (*args) that match the function signature + for i, arg in enumerate(args): + if i < len(parameters): + param = parameters[i] + expected_type = param.annotation + + if expected_type == inspect.Parameter.empty: + converted_args.append(arg) + else: + try: + # Use TypeAdapter to convert the argument + type_adapter = TypeAdapter(expected_type) + converted_args.append(type_adapter.validate_python(arg)) + except (ValueError, TypeError) as e: + print(f"WARNING - Could not convert arg[{i}]='{arg}' to {expected_type}: {e}") + converted_args.append(arg) + + # Include remaining *args as-is + if len(args) > len(parameters): + converted_args.extend(args[len(parameters):]) + + # Convert keyword arguments (*kwargs) that match the function signature + for name, param in sig.parameters.items(): + if param.kind in (inspect.Parameter.KEYWORD_ONLY, inspect.Parameter.POSITIONAL_OR_KEYWORD): + if name in kwargs: + expected_type = param.annotation + value = kwargs[name] + if expected_type != inspect.Parameter.empty: + try: + type_adapter = TypeAdapter(expected_type) + converted_kwargs[name] = type_adapter.validate_python(value) + except (ValueError, TypeError) as e: + print(f"WARNING - Could not convert kwarg '{name}'='{value}' to {expected_type}: {e}") + converted_kwargs[name] = value + else: + converted_kwargs[name] = value + + # Include any additional kwargs that weren't in the function signature + for k, v in kwargs.items(): + if k not in converted_kwargs: + converted_kwargs[k] = v + + # Ensure no positional argument conflicts with keyword arguments + for i, arg in enumerate(converted_args): + if i < len(parameters): + param_name = parameters[i].name + if param_name in converted_kwargs: + raise TypeError(f"Got multiple values for argument '{param_name}'") + + return tuple(converted_args), converted_kwargs + + +# Example function +def example_function( + local_specifier: str, + remote_specifier: str, + select: str = "url,domain_label,title,plain_text", + where: Optional[str] = "", + limit: Optional[int] = None, + files: str = "**/*.parquet", + explain: bool = False, + pq_batch_size: int = 1, + batch_size: int = 100, + prefetch: int = 2, + page_size: int = 10, + *args: Any, + **kwargs: Any +): + return locals() + + +# Example usage +args = ("example_local", "example_remote") +kwargs = { + "select": "url,domain_label,title", + "where": "url_suffix='at'", + "limit": "10", + "files": "**/*.parquet", + "explain": "false", + "pq_batch_size": "5", + "batch_size": "200", + "prefetch": "3", + "page_size": "15" +} + +converted_args, converted_kwargs = convert_args_kwargs_with_pydantic(example_function, args, kwargs) +print("Converted Args:", converted_args) +print("Converted Kwargs:", converted_kwargs) + +# Now you can call the function with the converted arguments +result = example_function(*converted_args, **converted_kwargs) +print("Function Result:", result) diff --git a/tests/test_cli.py b/tests/test_cli.py new file mode 100644 index 0000000..2fe4eec --- /dev/null +++ b/tests/test_cli.py @@ -0,0 +1,174 @@ +import shutil +import pytest +from click.testing import CliRunner +import owilix.cli +import os, fsspec +import re + + +def parse_summary_output(output: str) -> dict: + """ + Parses the summary output to extract key-value pairs where values are numbers, including totals. + + Args: + output (str): The output string from the CLI command. + + Returns: + dict: A dictionary with keys as the parsed keys and values as numerical values. + """ + # Dictionary to store the extracted key-value pairs + result = {} + + # Define regex patterns to match: + # - General key-value pairs where the value is a number + pattern = re.compile(r'^(.*?)\s+(\d[\d.,]*)\s*(?:\w*|)(?:\s*\w*|)(?:\s*GiB|)$') + # - Specific pattern for the "Total of N shown." line + total_pattern = re.compile(r'^Total of (\d+) shown\.$') + + # Iterate over each line in the output + for line in output.splitlines(): + line = line.strip() # Remove leading and trailing whitespace + + # Check if the line matches the general key-value pattern + match = pattern.match(line) + if match: + key = match.group(1).strip() + # Remove commas and convert to appropriate numerical type + value_str = match.group(2).replace(',', '') + if '.' in value_str: + value = float(value_str) + else: + value = int(value_str) + result[key] = value + continue + + # Check if the line matches the total pattern + total_match = total_pattern.match(line) + if total_match: + total_value = int(total_match.group(1)) + result["Total"] = total_value + + return result + +_fs = fsspec.filesystem("file") + +class TestCLI: + @pytest.fixture(scope="class") + def runner(self): + """Fixture that provides a Click CLI runner.""" + return CliRunner() + + @pytest.fixture(scope="session") + def temp_dir(self, tmpdir_factory, request): + """Fixture to provide a unique temporary directory for the session and ensure cleanup.""" + temp_dir = tmpdir_factory.mktemp("cli_tests", numbered=True) + + def cleanup(): + """Cleanup function to remove the temporary directory after the test session.""" + shutil.rmtree(str(temp_dir)) + + request.addfinalizer(cleanup) + return temp_dir + + def setup_method(self): + """Method to setup each test by registering commands.""" + owilix.cli.register_commands(owilix.cli.cli) + + def test_help_command(self, runner): + """Test the help command to ensure it displays the correct output.""" + result = runner.invoke(owilix.cli.cli, ['--help'], prog_name='owilix.cli') + assert result.exit_code == 0 + assert "Main command line interface group for OWI management tools." in result.output + + def test_local_empty_ls(self, runner, temp_dir): + """Test the 'local' command with a subcommand and specifier.""" + result = runner.invoke(owilix.cli.cli, ['--target', str(temp_dir), 'local', 'ls', 'all'], prog_name='owilix.cli') + assert result.exit_code == 0 + assert len(_fs.ls(str(temp_dir))) == 1 + assert _fs.exists(os.path.join(str(temp_dir),".logs")) + assert _fs.exists(os.path.join(str(temp_dir), ".logs","events.json")) + assert result.output=='Fetching datasets for specifier all\nNo data available to display.\n' + + @pytest.mark.parametrize("specifier", ["all", "lrz:2023-10-31", "lrz:latest", "lrz:2023-10-31#7", "lrz:2023-10-31/collectionName=main;resourceType=owi"]) + def test_remote_ls_non_empty_lrz(self, runner, temp_dir, specifier): + """Test the 'local' command with a subcommand and specifier.""" + result = runner.invoke(owilix.cli.cli, ['--target', str(temp_dir), 'remote', 'ls', specifier], prog_name='owilix.cli') + assert result.exit_code == 0 + summary = parse_summary_output(result.output) + assert summary["Total Files"]>0 + assert summary["DataCenter lrz"] > 0 + assert summary["resourceType owi"] > 0 + assert summary["Total"] > 0 + + @pytest.mark.parametrize("specifier", ["all", "it4i:2023-12-03", "it4i:latest", "it4i:2023-12-31#7", + "it4i:2023-12-31#7/collectionName=main;resourceType=owi"]) + def test_remote_ls_non_empty_it4i(self, runner, temp_dir, specifier): + """Test the 'local' command with a subcommand and specifier.""" + result = runner.invoke(owilix.cli.cli, ['--target', str(temp_dir), 'remote', 'ls', specifier], + prog_name='owilix.cli') + assert result.exit_code == 0 + summary = parse_summary_output(result.output) + assert summary["Total Files"] > 0 + assert summary["DataCenter it4i"] > 0 + assert summary["resourceType owi"] > 0 + assert summary["Total"] > 0 + + def test_remote_pull_files(self, runner, temp_dir): + """Test the 'remote' command with a subcommand and specifier.""" + result = runner.invoke(owilix.cli.cli, ['--yes', '--target', str(temp_dir), 'remote', 'pull', + 'it4i:2023-12-03', 'files=**/language=slv/*'], prog_name='owilix.cli') + assert result.exit_code == 0 + assert "Fetching files for OWI-Open Web Index-main.owi@it4i-2023-12-3:2023-12-3" in result.output + assert "Found 2 remote files. Syncing with local files" in result.output # Adjust this check based on expected output + result = runner.invoke(owilix.cli.cli, ['--target', str(temp_dir), 'local', 'ls', 'all'], + prog_name='owilix.cli') + assert result.exit_code == 0 + assert len(_fs.ls(str(temp_dir))) == 2 + assert _fs.exists(os.path.join(str(temp_dir), "public")) + assert any([f for f in _fs.ls(os.path.join(str(temp_dir),"public","main")) if f.endswith("json")]) + + + def test_config_command(self, runner, temp_dir): + result = runner.invoke(owilix.cli.cli, ['config', 'set', 'showfields=url,domain_label', '--target', str(temp_dir)], prog_name='owilix.cli') + assert result.exit_code == 0 + assert "config" in result.output # Adjust this check based on expected output + # Verify that configuration was correctly set in the target directory + config_path = os.path.join(temp_dir, "owilix.cfg") + assert os.path.exists(config_path) + + def test_query_command(self, runner, temp_dir): + """Test the 'query' command with a subcommand and options.""" + result = runner.invoke(owilix.cli.cli, ['query', 'run', '--local', 'all', '--remote', 'lrz', '--target', str(temp_dir)], prog_name='owilix.cli') + assert result.exit_code == 0 + assert "query" in result.output # Adjust this check based on expected output + + def test_clean_command(self, runner, temp_dir): + """Test the 'clean' command.""" + # Pre-create a '.env' file in the temp_dir + env_path = os.path.join(temp_dir, '.env') + with open(env_path, 'w') as f: + f.write('test content') + + result = runner.invoke(owilix.cli.cli, ['clean', '--target', str(temp_dir)], prog_name='owilix.cli') + assert result.exit_code == 0 + assert "Cleaning" in result.output # Check for expected 'clean' command output + + # Verify that the '.env' file has been removed + assert not os.path.exists(env_path) + + def test_logs_command(self, runner, temp_dir): + """Test the 'logs' command with a module argument.""" + result = runner.invoke(owilix.cli.cli, ['--target', str(temp_dir), 'logs', 'lexis', ], prog_name='owilix.cli') + assert result.exit_code == 0 + result = runner.invoke(owilix.cli.cli, ['--target', str(temp_dir), 'logs', 'errors', ], prog_name='owilix.cli') + assert result.exit_code == 0 + assert "log" in result.output + result = runner.invoke(owilix.cli.cli, ['--target', str(temp_dir), 'logs', 'events', ], prog_name='owilix.cli') + assert result.exit_code == 0 + assert "[" in result.output and "]" in result.output + + def test_invalid_command(self, runner): + """Test an invalid command to ensure proper error handling.""" + result = runner.invoke(owilix.cli.cli, ['nonexistent'], prog_name='owilix.cli') + assert result.exit_code != 0 + assert "No such command" in result.output