aboutsummaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorMax Nanis2026-08-29 22:00:43 -0700
committerMax Nanis2026-08-29 22:00:43 -0700
commitdda0067fdfe563e8270562f2e0c9da7c74c37a0b (patch)
tree53e742f9c7b585b1a8e46e4fe4f1fd04c629d196
parent8219bf4814a7aa374a500516f2254f39085f0357 (diff)
downloadgeneralresearch-dda0067fdfe563e8270562f2e0c9da7c74c37a0b.tar.gz
generalresearch-dda0067fdfe563e8270562f2e0c9da7c74c37a0b.zip
pylint ⭕️ cyclic-import checks
-rw-r--r--generalresearch/config.py5
-rw-r--r--generalresearch/grliq/managers/__init__.py34
-rw-r--r--generalresearch/grliq/models/forensic_summary.py6
-rw-r--r--generalresearch/incite/__init__.py4
-rw-r--r--generalresearch/incite/base.py10
-rw-r--r--generalresearch/incite/collections/__init__.py692
-rw-r--r--generalresearch/incite/collections/base.py693
-rw-r--r--generalresearch/incite/collections/thl_marketplaces.py2
-rw-r--r--generalresearch/incite/defaults.py10
-rw-r--r--generalresearch/incite/mergers/__init__.py301
-rw-r--r--generalresearch/incite/mergers/base.py301
-rw-r--r--generalresearch/incite/mergers/foundations/enriched_session.py2
-rw-r--r--generalresearch/incite/mergers/foundations/enriched_task_adjust.py2
-rw-r--r--generalresearch/incite/mergers/foundations/enriched_wall.py4
-rw-r--r--generalresearch/incite/mergers/foundations/user_id_product.py2
-rw-r--r--generalresearch/incite/mergers/pop_ledger.py2
-rw-r--r--generalresearch/incite/mergers/ym_survey_wall.py2
-rw-r--r--generalresearch/incite/mergers/ym_wall_summary.py2
-rw-r--r--generalresearch/incite/schemas/thl_web.py4
-rw-r--r--generalresearch/models/network/label.py1
-rw-r--r--generalresearch/models/thl/definitions.py2
-rw-r--r--generalresearch/models/thl/ipinfo.py22
-rw-r--r--generalresearch/models/thl/ledger.py63
-rw-r--r--generalresearch/models/thl/ledger_example.py64
-rw-r--r--generalresearch/models/thl/maxmind/__init__.py0
-rw-r--r--generalresearch/models/thl/maxmind/definitions.py22
-rw-r--r--generalresearch/models/thl/session.py25
-rw-r--r--generalresearch/models/thl/user_iphistory.py9
-rw-r--r--pyproject.toml6
-rw-r--r--test_utils/grliq/conftest.py50
-rw-r--r--test_utils/incite/mergers/conftest.py2
31 files changed, 1161 insertions, 1183 deletions
diff --git a/generalresearch/config.py b/generalresearch/config.py
index c414069..73af565 100644
--- a/generalresearch/config.py
+++ b/generalresearch/config.py
@@ -115,9 +115,8 @@ class GRLBaseSettings(BaseSettings):
amt_bonus_cashout_method_id: str | None = Field(default=None)
amt_assignment_cashout_method_id: str | None = Field(default=None)
- # --- Maxmind Configuration ---
- maxmind_account_id: str | None = Field(default=None)
- maxmind_license_key: str | None = Field(default=None)
+ # --- GRIP Configuration ---
+ grip_token: str | None = Field(default=None)
EXAMPLE_PRODUCT_ID = "1108d053e4fa47c5b0dbdcd03a7981e7"
diff --git a/generalresearch/grliq/managers/__init__.py b/generalresearch/grliq/managers/__init__.py
index 849b6c2..e69de29 100644
--- a/generalresearch/grliq/managers/__init__.py
+++ b/generalresearch/grliq/managers/__init__.py
@@ -1,34 +0,0 @@
-from generalresearch.grliq.models.forensic_data import GrlIqData
-from generalresearch.grliq.models.forensic_result import (
- GrlIqCheckerResults,
- GrlIqForensicCategoryResult,
-)
-
-DUMMY_GRLIQ_DATA = [
- {
- "data": GrlIqData.model_validate_json(
- """{"mid": "3722ed29314940fabd37b42d808dcf5a", "uuid": "b11441da5a854dfbb8401d4c32e56db5", "phase": "offerwall-enter", "events": null, "vendor": "Google Inc.", "app_name": "Netscape", "calendar": "gregory", "language": "en-US", "platform": "Linux x86_64", "timezone": "America/Mexico_City", "client_ip": "131.196.250.250", "timestamp": "2025-02-27T16:05:34-06:00", "webrtc_ip": "131.196.250.250", "created_at": "2025-02-27T22:05:35.370589Z", "language_2": "en-US", "language_3": null, "platform_2": "Linux x86_64", "platform_3": null, "prefetched": true, "product_id": "d0606a0b5d034a8d81b1e3579d1f76fd", "webgl_flag": true, "webgl_hash": "da27e1b9b660057a3f5e185d3f5deabe", "canvas_hash": "14ed764326ec454d976c322261d99f16", "color_gamut": "3", "country_iso": "mx", "inner_width": 612, "outer_width": 1813, "product_sub": "20030107", "audio_codecs": "1,1,1,1,1,3,1,3,1,3,3,1,1,3,3,3,3,1,3,3,3,2,1,1", "cookie_check": "", "graphics_api": "WebKit WebGL", "inner_height": 1174, "mouse_events": null, "ontouchstart": false, "outer_height": 1261, "plugins_hash": "4c05fa2f766a444d4f253ead792c8b0e|2", "screen_width": 2560, "video_codecs": "1,3,3,3,3,3,3,3,3,3,1,1,1,1,1,1,3,1,1,1,3,3,1", "webgl_hash_2": "fc73fd5db75e2c36222fe34251be3971", "webrtc_error": false, "window_opera": false, "battery_level": 0.9, "canvas_hash_2": "bd11ebbf5c26fd20e0217820b4159752", "dynamic_range": false, "error_message": "Cannot read", "forced_colors": false, "math_result_1": "1.9275814160560204e-50", "math_result_2": "1.6182817135715877", "screen_height": 1440, "webgl_check_1": true, "webgl_context": "webgl2", "window_chrome": true, "connection_rtt": 150, "history_length": 16, "user_agent_str": "Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/131.0.0.0 Safari/537.36", "web_sql_exists": false, "calender_locale": "en-US", "connection_type": "", "inverted_colors": true, "navigator_brave": false, "product_user_id": "d1d55df1-959e-4740-b77c-fa1f4fc457ae", "request_headers": {"host": "test", "accept": "*/*", "connection": "keep-alive", "user-agent": "python-httpx/0.27.0", "content-length": "3646", "accept-encoding": "gzip, deflate", "x-forwarded-for": "131.196.250.250"}, "timezone_offset": 360, "webrtc_local_ip": "50486637-6b64-4812-b10a-0a75337c31bd.local", "battery_charging": true, "client_ip_detail": {"continent_code": "EU", "continent_name": "Europe", "country_name": "France", "is_in_european_union": true, "ip": "131.196.250.250", "isp": null, "latitude": null, "city_name": null, "longitude": null, "time_zone": null, "user_type": null, "country_iso": "mx", "postal_code": null, "is_anonymous": null, "accuracy_radius": null, "static_ip_score": null, "subdivision_1_iso": null, "subdivision_2_iso": null, "subdivision_1_name": null, "subdivision_2_name": null, "registered_country_iso": null}, "max_touch_points": 0, "numbering_system": "latn", "path_fingerprint": 3252, "prefers_contrast": "0", "rendering_engine": "WebKit", "timezone_success": "pass", "user_agent_hints": {"model": null, "brands": [{"brand": "Google Chrome", "version": "131"}, {"brand": "Chromium", "version": "131"}, {"brand": "Not_A Brand", "version": "24"}], "mobile": false, "bitness": "64", "platform": "Linux", "brands_full": [{"brand": "Google Chrome", "version": "131.0.6778.204"}, {"brand": "Chromium", "version": "131.0.6778.204"}, {"brand": "Not_A Brand", "version": "24.0.0.0"}], "architecture": "x86", "platform_version": "6.2.0"}, "user_agent_str_2": null, "webgl_extensions": "EXT_clip_control|EXT_color_buffer_float|EXT_color_buffer_half_float|EXT_conservative_depth|EXT_depth_clamp|EXT_disjoint_timer_query_webgl2|EXT_float_blend|EXT_polygon_offset_clamp|EXT_render_snorm|EXT_texture_compression_bptc|EXT_texture_compression_rgtc|EXT_texture_filter_anisotropic|EXT_texture_mirror_clamp_to_edge|EXT_texture_norm16|KHR_parallel_shader_compile|NV_shader_noperspective_interpolation|OES_draw_buffers_indexed|OES_sample_variables|OES_shader_multisample_interpolation|OES_texture_float_linear|OVR_multiview2|WEBGL_blend_func_extended|WEBGL_clip_cull_distance|WEBGL_compressed_texture_astc|WEBGL_compressed_texture_etc|WEBGL_compressed_texture_etc1|WEBGL_compressed_texture_s3tc|WEBGL_compressed_texture_s3tc_srgb|WEBGL_debug_renderer_info|WEBGL_debug_shaders|WEBGL_lose_context|WEBGL_multi_draw|WEBGL_polygon_mode|WEBGL_provoking_vertex|WEBGL_stencil_texturing", "webrtc_ip_detail": {"continent_code": "EU", "continent_name": "Europe", "country_name": "France", "is_in_european_union": true, "ip": "131.196.250.250", "isp": null, "latitude": null, "city_name": null, "longitude": null, "time_zone": null, "user_type": null, "country_iso": "mx", "postal_code": null, "is_anonymous": null, "accuracy_radius": null, "static_ip_score": null, "subdivision_1_iso": null, "subdivision_2_iso": null, "subdivision_1_name": null, "subdivision_2_name": null, "registered_country_iso": null}, "chrome_extensions": "", "execution_time_ms": 371.0999999642372, "graphics_renderer": "WebGL 2.0 (OpenGL ES 3.0 Chromium)", "keyboard_detected": true, "mime_types_length": 2, "request_fs_exists": true, "audio_context_flag": "pass", "audio_context_hash": "9307303774dec3248c18a939392090da", "canvas_fingerprint": 258, "canvas_pixel_check": false, "device_pixel_ratio": 1.0, "indexedDbData_blob": true, "navigator_keys_len": 79, "no_edge_pdf_plugin": false, "screen_avail_width": 2560, "webdriver_detected": false, "window_orientation": 0, "connection_downlink": 10.0, "navigator_webdriver": false, "non_native_function": false, "screen_avail_height": 1400, "supported_fonts_str": "72|768|262144|1073741824|0|0|540672|73728|7340032|1342177280|117446656|256|16|0|543|4290797636|1677723648|4168998400|0|1048576|262144|268500994|1342177280|262144|125829376|37888000|0|435363842|0|2147483648|109543424|1880099872|268435471", "text_2d_fingerprint": "bfcce91c9e71d11af7b14dbee4c75f83", "webrtc_is_supported": "pass", "canvas_support_level": "full", "do_not_track_enabled": "1", "hardware_concurrency": 12, "keyboard_layout_size": 48, "prefers_color_scheme": false, "webgl_max_anisotropy": 16, "battery_charging_time": 0.0, "browser_by_properties": "c", "eval_to_string_length": 33, "performance_loop_time": 0.09999996423721313, "session_storage_check": "pass", "unmasked_vendor_webgl": "Google Inc. (Intel)", "hardware_concurrency_2": 12, "hardware_concurrency_3": null, "localStorage_available": true, "memory_jsHeapSizeLimit": 4294705152, "mozilla_web_app_exists": false, "navigator_deviceMemory": 8.0, "navigator_java_enabled": false, "prefers_reduced_motion": false, "storage_estimate_quota": 1178717110272, "webdriver_detected_msg": "", "window_active_x_object": false, "window_external_exists": true, "color_depth_pixel_depth": "24-24", "indexedDbData_available": true, "navigator_cookieEnabled": true, "unmasked_renderer_webgl": "ANGLE (Intel, Mesa Intel(R) Graphics (RPL-P), OpenGL 4.6)", "battery_discharging_time": 0.0, "connection_effectiveType": "4g", "non_native_function_flag": "", "speech_synthesis_voice_1": "Google Bahasa Indonesia", "window_client_information": true, "audio_compressor_reduction": 20.538288116455078, "navigator_mediaDevices_len": 3, "audio_intensity_fingerprint": 124.04347527516074, "speech_synthesis_voice_hash": "8010ee3313813de521e48e63bd5a6f13", "microsoft_credentials_exists": false, "window_installTrigger_exists": false, "speech_synthesis_voices_count": 19, "webgl_shading_language_version": "WebGL GLSL ES 3.00 (OpenGL ES GLSL ES 3.0 Chromium)", "error_message_stack_access_count": 0, "speech_synthesis_avail_voices_count": 19, "error_message_stack_access_count_worker": 0}"""
- ),
- "result_data": GrlIqCheckerResults.model_validate_json(
- """{"uuid": "b11441da5a854dfbb8401d4c32e56db5", "check_codecs": {"score": 0}, "check_timezone": {"score": 0}, "check_timestamp": {"score": 0}, "check_user_type": {"score": 0}, "check_ip_changes": {"score": 0}, "check_ip_country": {"score": 0}, "check_environment": {"score": 0}, "check_ip_timezone": {"score": 0}, "check_isp_changes": {"score": 0}, "check_useragent_js": {"score": 0}, "check_required_fonts": {"score": 0}, "check_user_anonymous": {"score": 0}, "check_webrtc_success": {"score": 0}, "check_seen_timestamps": {"msg": "duplicate timestamp", "score": 100}, "check_country_timezone": {"score": 0}, "check_prohibited_fonts": {"score": 0}, "check_timezone_changes": {"score": 0}, "check_execution_time_ms": {"msg": "duplicate execution_time_ms", "score": 100}, "check_fingerprint_reuse": {"score": 0}, "check_fingerprint_cycling": {"score": 0}, "check_ip_webrtc_ip_detail": {"score": 0}, "check_environment_critical": {"score": 0}, "check_useragent_other_enums": {"score": 0}, "check_useragent_ip_properties": {"score": 0}, "check_useragent_data_properties": {"score": 0}, "check_useragent_device_family_brand": {"score": 0}}"""
- ),
- "category_result": GrlIqForensicCategoryResult.model_validate_json(
- """{"uuid": "b11441da5a854dfbb8401d4c32e56db5", "is_bot": 0, "is_tampered": 100, "is_velocity": 0, "is_anonymous": 0, "suspicious_ip": 0, "is_oscillating": 0, "is_teleporting": 0, "is_inconsistent": 0, "platform_ip_inconsistent": 0}"""
- ),
- "fraud_score": 100,
- "is_attempt_allowed": False,
- },
- {
- "data": GrlIqData.model_validate_json(
- """{"mid": "35f6f5c30bc74ea7ac4aca7b40a02352", "uuid": "d54509f2f310499f8ab74839b10b2a41", "phase": "offerwall-enter", "events": null, "vendor": "Google Inc.", "app_name": "Netscape", "calendar": "gregory", "language": "en-US", "platform": "Linux x86_64", "timezone": "America/Los_Angeles", "client_ip": "104.9.125.144", "timestamp": "2025-02-28T11:34:39-08:00", "webrtc_ip": "172.56.209.195", "created_at": "2025-02-28T19:34:39.681872Z", "language_2": "en-US", "language_3": null, "platform_2": "Linux x86_64", "platform_3": null, "prefetched": true, "product_id": "d0606a0b5d034a8d81b1e3579d1f76fd", "webgl_flag": true, "webgl_hash": "da27e1b9b660057a3f5e185d3f5deabe", "canvas_hash": "e6e4d17da26050ce85ad00d3c6ea999e", "color_gamut": "3", "country_iso": "us", "inner_width": 841, "outer_width": 1680, "product_sub": "20030107", "audio_codecs": "1,1,1,1,1,3,1,3,1,3,3,1,1,3,3,3,3,1,3,3,3,2,1,1", "cookie_check": "", "graphics_api": "WebKit WebGL", "inner_height": 891, "mouse_events": null, "ontouchstart": false, "outer_height": 978, "plugins_hash": "4c05fa2f766a444d4f253ead792c8b0e|2", "screen_width": 1680, "video_codecs": "1,3,3,3,3,3,3,3,3,3,1,1,1,1,1,1,3,1,1,1,3,3,1", "webgl_hash_2": "fc73fd5db75e2c36222fe34251be3971", "webrtc_error": false, "window_opera": false, "battery_level": 0.41, "canvas_hash_2": "e0559d49b1864985cafc0d1c3a6b053c", "dynamic_range": false, "error_message": "Cannot read", "forced_colors": false, "math_result_1": "1.9275814160560204e-50", "math_result_2": "1.6182817135715877", "screen_height": 1050, "webgl_check_1": true, "webgl_context": "webgl2", "window_chrome": true, "connection_rtt": 100, "history_length": 11, "user_agent_str": "Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/131.0.0.0 Safari/537.36", "web_sql_exists": false, "calender_locale": "en-US", "connection_type": "", "inverted_colors": true, "navigator_brave": false, "product_user_id": "test-unit", "request_headers": {"dnt": "1", "host": "127.0.0.1:8081", "accept": "application/json, lk/null q=0.1", "origin": "http://127.0.0.1:8080", "referer": "http://127.0.0.1:8080/", "sec-ch-ua": "\\"Google Chrome\\";v=\\"131\\", \\"Chromium\\";v=\\"131\\", \\"Not_A Brand\\";v=\\"24\\"", "connection": "keep-alive", "user-agent": "Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/131.0.0.0 Safari/537.36", "content-type": "application/json", "content-length": "3313", "sec-fetch-dest": "empty", "sec-fetch-mode": "cors", "sec-fetch-site": "same-site", "accept-encoding": "gzip, deflate, br, zstd", "accept-language": "en-US,en;q=0.9", "sec-ch-ua-mobile": "?0", "sec-ch-ua-platform": "\\"Linux\\""}, "timezone_offset": 480, "webrtc_local_ip": "10.253.217.45,[2607:fb91:20c5:c6af:cda0:10b4:830a:a85e]", "battery_charging": false, "client_ip_detail": {"continent_code": "EU", "continent_name": "Europe", "country_name": "France", "is_in_european_union": true, "ip": "104.9.125.144", "isp": "AT&T Internet", "latitude": 37.3897, "city_name": "Mountain View", "longitude": -122.083, "time_zone": "America/Los_Angeles", "user_type": "residential", "country_iso": "us", "postal_code": "94041", "is_anonymous": false, "accuracy_radius": 5, "static_ip_score": 40.3, "subdivision_1_iso": "CA", "subdivision_2_iso": null, "subdivision_1_name": "California", "subdivision_2_name": null, "registered_country_iso": "us"}, "max_touch_points": 0, "numbering_system": "latn", "path_fingerprint": 3252, "prefers_contrast": "0", "rendering_engine": "WebKit", "timezone_success": "pass", "user_agent_hints": {"model": null, "brands": [{"brand": "Google Chrome", "version": "131"}, {"brand": "Chromium", "version": "131"}, {"brand": "Not_A Brand", "version": "24"}], "mobile": false, "bitness": "64", "platform": "Linux", "brands_full": [{"brand": "Google Chrome", "version": "131.0.6778.204"}, {"brand": "Chromium", "version": "131.0.6778.204"}, {"brand": "Not_A Brand", "version": "24.0.0.0"}], "architecture": "x86", "platform_version": "6.2.0"}, "user_agent_str_2": null, "webgl_extensions": "EXT_clip_control|EXT_color_buffer_float|EXT_color_buffer_half_float|EXT_conservative_depth|EXT_depth_clamp|EXT_disjoint_timer_query_webgl2|EXT_float_blend|EXT_polygon_offset_clamp|EXT_render_snorm|EXT_texture_compression_bptc|EXT_texture_compression_rgtc|EXT_texture_filter_anisotropic|EXT_texture_mirror_clamp_to_edge|EXT_texture_norm16|KHR_parallel_shader_compile|NV_shader_noperspective_interpolation|OES_draw_buffers_indexed|OES_sample_variables|OES_shader_multisample_interpolation|OES_texture_float_linear|OVR_multiview2|WEBGL_blend_func_extended|WEBGL_clip_cull_distance|WEBGL_compressed_texture_astc|WEBGL_compressed_texture_etc|WEBGL_compressed_texture_etc1|WEBGL_compressed_texture_s3tc|WEBGL_compressed_texture_s3tc_srgb|WEBGL_debug_renderer_info|WEBGL_debug_shaders|WEBGL_lose_context|WEBGL_multi_draw|WEBGL_polygon_mode|WEBGL_provoking_vertex|WEBGL_stencil_texturing", "webrtc_ip_detail": {"continent_code": "EU", "continent_name": "Europe", "country_name": "France", "is_in_european_union": true, "ip": "172.56.209.195", "isp": null, "latitude": null, "city_name": null, "longitude": null, "time_zone": null, "user_type": null, "country_iso": "us", "postal_code": null, "is_anonymous": null, "accuracy_radius": null, "static_ip_score": null, "subdivision_1_iso": null, "subdivision_2_iso": null, "subdivision_1_name": null, "subdivision_2_name": null, "registered_country_iso": null}, "chrome_extensions": "", "execution_time_ms": 924.5, "graphics_renderer": "WebGL 2.0 (OpenGL ES 3.0 Chromium)", "keyboard_detected": true, "mime_types_length": 2, "request_fs_exists": true, "audio_context_flag": "pass", "audio_context_hash": "9307303774dec3248c18a939392090da", "canvas_fingerprint": 258, "canvas_pixel_check": false, "device_pixel_ratio": 1.0, "indexedDbData_blob": true, "navigator_keys_len": 79, "no_edge_pdf_plugin": false, "screen_avail_width": 1680, "webdriver_detected": false, "window_orientation": 0, "connection_downlink": 10.0, "navigator_webdriver": false, "non_native_function": false, "screen_avail_height": 1010, "supported_fonts_str": "72|17152|327680|1073741824|0|0|540736|73728|7340032|1342177280|117446657|256|16|0|262687|4290797636|1677723648|4168998400|0|1048576|262144|268500994|1342177280|262144|125829376|37888000|0|435363842|0|2147483648|109543680|1880099888|301989903", "text_2d_fingerprint": "bfcce91c9e71d11af7b14dbee4c75f83", "webrtc_is_supported": "pass", "canvas_support_level": "full", "do_not_track_enabled": "1", "hardware_concurrency": 12, "keyboard_layout_size": 48, "prefers_color_scheme": false, "webgl_max_anisotropy": 16, "battery_charging_time": 0.0, "browser_by_properties": "c", "eval_to_string_length": 33, "performance_loop_time": 0.09999999962747097, "session_storage_check": "pass", "unmasked_vendor_webgl": "Google Inc. (Intel)", "hardware_concurrency_2": 12, "hardware_concurrency_3": null, "localStorage_available": true, "memory_jsHeapSizeLimit": 4294705152, "mozilla_web_app_exists": false, "navigator_deviceMemory": 8.0, "navigator_java_enabled": false, "prefers_reduced_motion": false, "storage_estimate_quota": 1178717110272, "webdriver_detected_msg": "", "window_active_x_object": false, "window_external_exists": true, "color_depth_pixel_depth": "24-24", "indexedDbData_available": true, "navigator_cookieEnabled": true, "unmasked_renderer_webgl": "ANGLE (Intel, Mesa Intel(R) Graphics (RPL-P), OpenGL 4.6)", "battery_discharging_time": 4844.0, "connection_effectiveType": "4g", "non_native_function_flag": "", "speech_synthesis_voice_1": "Google Bahasa Indonesia", "window_client_information": true, "audio_compressor_reduction": 20.538288116455078, "navigator_mediaDevices_len": 8, "audio_intensity_fingerprint": 124.04347527516074, "speech_synthesis_voice_hash": "8010ee3313813de521e48e63bd5a6f13", "microsoft_credentials_exists": false, "window_installTrigger_exists": false, "speech_synthesis_voices_count": 19, "webgl_shading_language_version": "WebGL GLSL ES 3.00 (OpenGL ES GLSL ES 3.0 Chromium)", "error_message_stack_access_count": 2, "speech_synthesis_avail_voices_count": 19, "error_message_stack_access_count_worker": 2}"""
- ),
- "result_data": GrlIqCheckerResults.model_validate_json(
- """{"uuid": "d54509f2f310499f8ab74839b10b2a41", "check_codecs": {"score": 0}, "check_timezone": {"score": 0}, "check_timestamp": {"score": 0}, "check_user_type": {"score": 0}, "check_ip_changes": {"score": 0}, "check_ip_country": {"score": 0}, "check_environment": {"msg": "error_message_stack_access_count: 2", "score": 100}, "check_ip_timezone": {"score": 0}, "check_isp_changes": {"score": 0}, "check_useragent_js": {"score": 0}, "check_required_fonts": {"score": 0}, "check_user_anonymous": {"score": 0}, "check_webrtc_success": {"score": 0}, "check_seen_timestamps": {"score": 0}, "check_country_timezone": {"score": 0}, "check_prohibited_fonts": {"score": 0}, "check_timezone_changes": {"score": 0}, "check_execution_time_ms": {"score": 0}, "check_fingerprint_reuse": {"score": 0}, "check_fingerprint_cycling": {"score": 0}, "check_ip_webrtc_ip_detail": {"score": 0}, "check_environment_critical": {"score": 0}, "check_useragent_other_enums": {"score": 0}, "check_useragent_ip_properties": {"score": 0}, "check_useragent_data_properties": {"score": 0}, "check_useragent_device_family_brand": {"score": 0}}"""
- ),
- "category_result": GrlIqForensicCategoryResult.model_validate_json(
- """{"uuid": "d54509f2f310499f8ab74839b10b2a41", "is_bot": 0, "is_tampered": 0, "is_velocity": 0, "is_anonymous": 0, "suspicious_ip": 0, "is_oscillating": 0, "is_teleporting": 0, "is_inconsistent": 10, "platform_ip_inconsistent": 0}"""
- ),
- "fraud_score": 10,
- "is_attempt_allowed": True,
- },
-]
diff --git a/generalresearch/grliq/models/forensic_summary.py b/generalresearch/grliq/models/forensic_summary.py
index 6b0e065..aaefdfb 100644
--- a/generalresearch/grliq/models/forensic_summary.py
+++ b/generalresearch/grliq/models/forensic_summary.py
@@ -10,6 +10,7 @@ from typing import (
)
import numpy as np
+from grip_client.enums import AccessType
from pydantic import (
BaseModel,
ConfigDict,
@@ -27,7 +28,6 @@ from generalresearch.grliq.models.forensic_result import (
)
from generalresearch.models.custom_types import AwareDatetimeISO, IPvAnyAddressStr
from generalresearch.models.thl.locales import CountryISO
-from generalresearch.models.thl.maxmind.definitions import UserType
example_rtt_percentiles = (
[133.332]
@@ -185,7 +185,9 @@ class IPTimingDataSummary(BaseModel):
client_ip: IPvAnyAddressStr = Field(examples=["123.123.123.123"])
country_iso: CountryISO = Field(examples=["us"])
server_location: Literal["fremont_ca"] = Field(default="fremont_ca")
- user_type: UserType | None = Field(default=None, examples=[UserType.RESIDENTIAL])
+ user_type: AccessType | None = Field(
+ default=None, examples=[AccessType.RESIDENTIAL]
+ )
expected_rtt_range: tuple[float, float] = Field(
description="The expected rtt range for this IP (based on country_iso/user_type) to server_location",
examples=[(45.193, 120.841)],
diff --git a/generalresearch/incite/__init__.py b/generalresearch/incite/__init__.py
index e69de29..8b60e4b 100644
--- a/generalresearch/incite/__init__.py
+++ b/generalresearch/incite/__init__.py
@@ -0,0 +1,4 @@
+import logging
+
+logging.basicConfig()
+LOG = logging.getLogger(f"{__name__}.incite")
diff --git a/generalresearch/incite/base.py b/generalresearch/incite/base.py
index 060df4e..9dfec4f 100644
--- a/generalresearch/incite/base.py
+++ b/generalresearch/incite/base.py
@@ -1,7 +1,6 @@
from __future__ import annotations
import glob
-import logging
import os
import re
import shutil
@@ -46,7 +45,8 @@ from pydantic.json_schema import SkipJsonSchema
from sentry_sdk import capture_exception
from generalresearch.config import is_debug
-from generalresearch.incite.collections import DFCollectionItem
+from generalresearch.incite import LOG
+from generalresearch.incite.collections.base import DFCollectionItem
from generalresearch.incite.schemas import (
ARCHIVE_AFTER,
empty_dataframe_from_schema,
@@ -54,16 +54,14 @@ from generalresearch.incite.schemas import (
from generalresearch.models.custom_types import AwareDatetimeISO
if TYPE_CHECKING:
- from generalresearch.incite.collections import DFCollection
+ from generalresearch.incite.collections.base import DFCollection
from generalresearch.incite.collections.thl_marketplaces import (
DFCollectionType,
)
- from generalresearch.incite.mergers import MergeCollection, MergeType
+ from generalresearch.incite.mergers.base import MergeCollection, MergeType
Collection = DFCollection | MergeCollection
-logging.basicConfig()
-LOG = logging.getLogger(f"{__name__}.incite")
# Item = Union["DFCollectionItem", "MergeCollectionItem"]
Item = Any
diff --git a/generalresearch/incite/collections/__init__.py b/generalresearch/incite/collections/__init__.py
index f1e1b77..e69de29 100644
--- a/generalresearch/incite/collections/__init__.py
+++ b/generalresearch/incite/collections/__init__.py
@@ -1,692 +0,0 @@
-from __future__ import annotations
-
-import os
-import subprocess
-import time
-from datetime import datetime
-from enum import StrEnum
-from sys import platform
-from typing import Any
-
-import dask
-import dask.dataframe as dd
-import pandas as pd
-import pyarrow as pa
-import pyarrow.parquet as pq
-from dask.distributed import Client as DaskClient
-from dask.distributed import Future
-from distributed import as_completed
-from more_itertools import chunked
-from pandera.pandas import DataFrameSchema
-from psycopg import Cursor
-from pydantic import Field, FilePath, ValidationInfo, field_validator
-from sentry_sdk import capture_exception
-
-from generalresearch.incite.base import LOG, CollectionBase, CollectionItemBase
-from generalresearch.incite.schemas import (
- ARCHIVE_AFTER,
- ORDER_KEY,
- PARTITION_ON,
- empty_dataframe_from_schema,
-)
-from generalresearch.incite.schemas.thl_marketplaces import (
- InnovateSurveyHistorySchema,
- MorningSurveyTimeseriesSchema,
- SagoSurveyHistorySchema,
- SpectrumSurveyTimeseriesSchema,
-)
-from generalresearch.incite.schemas.thl_web import (
- LedgerSchema,
- THLIPInfoSchema,
- THLSessionSchema,
- THLTaskAdjustmentSchema,
- THLUserSchema,
- THLWallSchema,
- TransactionMetadataColumns,
- TxMetaSchema,
- TxSchema,
- UserHealthAuditLogSchema,
- UserHealthIPHistorySchema,
- UserHealthIPHistoryWSSchema,
-)
-from generalresearch.pg_helper import PostgresConfig
-
-DT_STR = "%Y-%m-%d %H:%M:%S"
-
-
-class DFCollectionType(StrEnum):
- TEST = "test"
-
- USER = "thl_user"
- SESSION = "thl_session"
- WALL = "thl_wall"
- TASK_ADJUSTMENT = "thl_taskadjustment"
- IP_INFO = "thl_ipinformation"
-
- AUDIT_LOG = "userhealth_auditlog"
- IP_HISTORY = "userhealth_iphistory"
- IP_HISTORY_WS = "userhealth_iphistory_ws"
-
- LEDGER = "ledger"
-
- INNOVATE_SURVEY_HISTORY = "innovate_surveyhistory"
- MORNING_SURVEY_TIMESERIES = "morning_surveytimeseries"
- SAGO_SURVEY_HISTORY = "sago_surveyhistory"
- SPECTRUM_SURVEY_TIMESERIES = "spectrum_surveytimeseries"
-
-
-DFCollectionTypeSchemas = {
- DFCollectionType.USER: THLUserSchema,
- DFCollectionType.WALL: THLWallSchema,
- DFCollectionType.SESSION: THLSessionSchema,
- DFCollectionType.IP_INFO: THLIPInfoSchema,
- DFCollectionType.TASK_ADJUSTMENT: THLTaskAdjustmentSchema,
- DFCollectionType.IP_HISTORY: UserHealthIPHistorySchema,
- DFCollectionType.IP_HISTORY_WS: UserHealthIPHistoryWSSchema,
- DFCollectionType.AUDIT_LOG: UserHealthAuditLogSchema,
- DFCollectionType.LEDGER: LedgerSchema,
- DFCollectionType.INNOVATE_SURVEY_HISTORY: InnovateSurveyHistorySchema,
- DFCollectionType.MORNING_SURVEY_TIMESERIES: MorningSurveyTimeseriesSchema,
- DFCollectionType.SAGO_SURVEY_HISTORY: SagoSurveyHistorySchema,
- DFCollectionType.SPECTRUM_SURVEY_TIMESERIES: SpectrumSurveyTimeseriesSchema,
-}
-
-
-class DFCollectionItem(CollectionItemBase):
-
- # --- Properties ---
- @property
- def filename(self) -> str:
- return (
- f"{self._collection.data_type.name.lower()}-{self._collection.offset}"
- f"-{self.start.strftime('%Y-%m-%d-%H-%M-%S')}.parquet"
- )
-
- # --- Methods ---
-
- def has_postgres(self) -> bool:
- if self._collection.pg_config is None:
- return False
-
- connected = True
- try:
- self._collection.pg_config.execute_sql_query("""SELECT 1;""")
- except AssertionError:
- connected = False
-
- return connected
-
- def has_db(self) -> bool:
- return self.has_mysql() or self.has_postgres()
-
- def update_partial_archive(self) -> bool:
- if not self.valid_archive(self.partial_path, sample=1000):
- LOG.error(f"invalid partial archive: {self.partial_path}")
- return self.create_partial_archive()
- df = pq.ParquetDataset(self.partial_path).read().to_pandas()
-
- order_key = self._collection._schema.metadata[ORDER_KEY]
- archive_after = self._collection._schema.metadata[ARCHIVE_AFTER]
-
- partial_max = df[order_key].max().to_pydatetime()
-
- since = partial_max - archive_after
- since = max([since, self.start]) # don't allow to query before the item's start
- df = df[df[order_key] < since].copy()
-
- _df = self.from_db(since=since)
-
- if _df is not None:
- df = pd.concat([df, _df])
- self.to_archive(ddf=dd.from_pandas(df, npartitions=1), is_partial=True)
- else:
- # The update to the partial returned no rows, but the partial
- # still exists, so we'll continue with whatever was calling this.
- # We don't need to re-write the partial or really do anything.
- pass
- return True
-
- def create_partial_archive(self) -> bool:
- _df = self.from_db()
- if _df is None:
- # Returned no rows, but the period is not closed, so we
- # don't want to mark as empty. Do nothing.
- return False
- return self.to_archive(ddf=dd.from_pandas(_df, npartitions=1), is_partial=True)
-
- # --- ORM / Data handlers---
- def to_dict(self) -> dict[str, Any]:
- return self._to_dict()
-
- def from_db(self, since: datetime | None = None) -> pd.DataFrame | None:
- if self._collection.data_type == DFCollectionType.LEDGER:
- assert since is None, "Shouldn't pass since for Ledger item"
- assert self._collection.pg_config is not None
- return self.from_postgres_ledger()
- else:
- return self.from_db_standard(since=since)
-
- def from_db_standard(self, since: datetime | None = None) -> pd.DataFrame | None:
- assert (
- self._collection.data_type != DFCollectionType.LEDGER
- ), "Can't call from_postgres_standard for Ledger DFCollectionItem"
-
- start, finish = self.start, self.finish
- LOG.debug(
- f"{self._collection.data_type.value}.from_postgres("
- f"start={start.strftime(DT_STR)}, "
- f"finish={finish.strftime(DT_STR)})"
- )
- coll = self._collection
- schema = coll._schema
- pg_config = coll.pg_config
- assert pg_config, "Must provide PostgresConfig"
-
- start = since or start
- order_key = schema.metadata[ORDER_KEY]
- cols = list(schema.columns.keys()) + [schema.index.name]
- cols_str = ", ".join(cols)
-
- try:
- res = pg_config.execute_sql_query(
- query=f"""
- SELECT {cols_str}
- FROM {coll.data_type.value}
- WHERE {order_key} >= %s AND {order_key} < %s;
- """,
- params=[start, finish],
- )
- except AssertionError as e:
- capture_exception(error=e)
- LOG.error(f"_from_postgres Exception: {e}")
- return None
-
- if not res:
- LOG.warning("_from_postgres query returned nothing")
- # Return an empty df.DataFrame with the correct columns
- return empty_dataframe_from_schema(coll._schema)
-
- df = pd.DataFrame.from_records(res).set_index(coll._schema.index.name)
- df = self.validate_df(df=df)
-
- if df is None:
- LOG.warning("_from_postgres query results failed validation")
- # Schema validation can fail...
- return None
-
- return df
-
- def from_postgres_ledger(self) -> pd.DataFrame | None:
- assert (
- self._collection.data_type == DFCollectionType.LEDGER
- ), "Can only call from_postgres_ledger on Ledger DFCollectionItem"
-
- start, finish = self.start, self.finish
- LOG.info(
- f"{self._collection.data_type.value}.from_postgres_ledger("
- f"start={start.strftime(DT_STR)}, "
- f"finish={finish.strftime(DT_STR)})"
- )
-
- coll = self._collection
- assert coll.pg_config, "Must provide PostgresConfig"
- pg_config: PostgresConfig = coll.pg_config
-
- limit = 20000
- offset = 0
- res = []
- while True:
- LOG.info(
- f"{self._collection.data_type.value}.from_postgres_ledger({limit=}, {offset=})"
- )
- chunk = pg_config.execute_sql_query(
- query=f"""
- SELECT lt.id AS tx_id, lt.created, lt.ext_description, lt.tag,
- le.id AS entry_id, le.direction, le.amount, le.account_id,
- la.display_name, la.qualified_name, la.account_type,
- la.normal_balance, la.reference_type, la.reference_uuid,
- la.currency
- FROM ledger_transaction AS lt
- LEFT JOIN ledger_entry AS le
- ON lt.id = le.transaction_id
- LEFT JOIN ledger_account AS la
- ON la.uuid = le.account_id
- WHERE lt.created >= %s AND lt.created < %s
- AND le.id IS NOT NULL
- ORDER BY lt.created
- LIMIT {limit} OFFSET {offset};
- """,
- params=[start, finish],
- )
- res.extend(chunk)
- if not chunk:
- break
- offset += limit
-
- if len(res) == 0:
- return None
-
- # Note (AND le.id IS NOT NULL): It is possible we have transactions with
- # no ledger entries. This is because the transaction creation failed
- # for some reason. The ledger is not unbalanced, it is just an orphan
- # transaction. Just skip those here.
-
- tx_df = TxSchema.validate(
- check_obj=pd.DataFrame.from_records(res).set_index("entry_id"),
- lazy=True,
- )
-
- tx_ids = list(tx_df["tx_id"].unique())
- metadata_res = []
- # "MySQL server has gone away" if this is too big
- conn = pg_config.make_connection()
- c: Cursor = conn.cursor()
- for chunk in chunked(tx_ids, n=5_000):
- c.execute(
- query="""
- SELECT ltm.transaction_id AS tx_id,
- ltm.id AS tx_metadata_id,
- ltm.key, ltm.value
- FROM ledger_transactionmetadata AS ltm
- WHERE ltm.transaction_id = ANY(%s);
- """,
- params=[chunk],
- )
- metadata_res += c.fetchall()
-
- conn.close()
-
- tx_meta = (
- pd.DataFrame(
- TxMetaSchema.validate(
- check_obj=pd.DataFrame.from_records(metadata_res).set_index(
- ["tx_id", "tx_metadata_id"]
- ),
- lazy=True,
- ).pivot(columns="key", values="value"),
- # This makes sure we expand to have all the possible columns
- columns=[e.value for e in TransactionMetadataColumns],
- )
- .groupby("tx_id")
- .first()
- )
-
- df = tx_df.merge(tx_meta, how="left", left_on="tx_id", right_index=True)
- df = self.validate_df(df=df)
-
- if df is None:
- # Schema validation can fail...
- return None
-
- return df
-
- def to_archive(
- self,
- ddf: dd.DataFrame,
- is_partial: bool = False,
- overwrite: bool = False,
- ) -> bool:
- """
- :returns: bool (saved_successful)
- """
- assert isinstance(ddf, dd.DataFrame), "must pass dask df"
-
- client: DaskClient | None = self._collection._client
- # client = None
-
- if client:
- row_len = client.compute(collections=ddf.shape[0], sync=True)
- else:
- row_len = len(ddf.index)
- is_empty = row_len == 0
-
- if is_partial:
- return self.to_archive_numbered_partial(ddf=ddf)
- else:
- return self._to_archive(
- ddf=ddf,
- is_empty=is_empty,
- overwrite=overwrite,
- )
-
- def _to_archive(
- self,
- ddf: dd.DataFrame | None,
- is_empty: bool,
- overwrite: bool = False,
- ) -> bool:
- """
- For archiving an item. Will write an empty file if ddf is empty. This
- is NOT for writing partials.
-
- :returns: bool (saved_successful)
- """
-
- if ddf is None:
- return False
-
- should_archive = self.should_archive()
- if not should_archive:
- LOG.warning(f"Cannot create archive for such new data: {self.path}")
- return False
-
- if overwrite is False:
- has_archive = self.has_archive(include_empty=True)
- if has_archive:
- LOG.warning(f"archive already exists: {self.path}")
- return False
-
- if is_empty:
- # Create an .empty only if the Item is "archiveable" (which we checked above)
- self.set_empty()
- return True
-
- # Incase the file saving is interrupted, or otherwise fails
- # save it to a tmp file first, then rename once we can confirm
- # that it successfully loads
- tmp_path = self.tmp_path()
- try:
- schema = self._collection._schema
- assert schema
- assert schema.metadata
- partition = schema.metadata.get(PARTITION_ON)
-
- ddf.to_parquet(
- path=tmp_path,
- partition_on=partition,
- engine="pyarrow",
- overwrite=True,
- write_metadata_file=True,
- compression="brotli",
- )
-
- except (pa.ArrowInvalid, pa.ArrowIOError, OSError) as e:
- LOG.exception(e)
- self.delete_archive(tmp_path)
- return False
-
- # It was saved, but the file seems to be corrupt
- if not self.valid_archive(tmp_path):
- LOG.error(f"not valid archive: {tmp_path}")
- self.delete_archive(tmp_path)
- # File did not save correctly so return it as saved=False
- return False
-
- # To debug, just set this key to auto expire in 5 seconds
- # RC.set(name=f"_to_archive:{self.path.as_posix()}", value=1, ex=15)
- # with RC.lock(f"_to_archive:{self.path.as_posix()}:lock", timeout=15):
-
- if os.path.isfile(tmp_path):
- # If the file was saved okay, seems okay, rename it
- os.replace(tmp_path, self.path)
- os.remove(tmp_path)
-
- if os.path.isdir(tmp_path):
- if os.path.exists(self.path.as_posix()):
- if overwrite:
- subprocess.call(["rm", "-r", self.path.as_posix()])
- time.sleep(1)
- else:
- LOG.error(f"already exists: {self.path.as_posix()}")
- return False
-
- if platform == "darwin":
- subprocess.call(["mv", tmp_path.as_posix(), self.path.as_posix()])
- else:
- # -T will (should) cause the mv to fail if path wasn't successfully deleted
- subprocess.call(["mv", "-T", tmp_path.as_posix(), self.path.as_posix()])
- return True
-
- def to_archive_numbered_partial(self, ddf: dd.DataFrame | None = None) -> bool:
- """
- For partial files/dirs only. Writes the .partial file with a number
- at the end (.partial.####) and then creates a symlink
- from .partial -> .partial.####
-
- :returns: bool (saved_successful)
- """
- if ddf is None:
- return False
-
- collection = self._collection
- schema = collection._schema
- assert schema
- client: DaskClient | None = collection._client
-
- next_numbered_path = self.next_numbered_path(self.partial_path)
- partial_path = self.partial_path
- # finish = self.finish
-
- # Make sure these are in the same dir. b/c the symlink has to be
- # relative, not an absolute path
- assert (
- partial_path.parent == next_numbered_path.parent
- ), "Can't have numbered_path in a different directory"
- target = (
- next_numbered_path.name
- ) # this is the symlink's target. it is a relative path (only the name)
-
- should_archive = self.should_archive()
- assert should_archive is False, "Don't write partial if the item is archiveable"
-
- if client:
- row_len = client.compute(collections=ddf.shape[0], sync=True)
- else:
- row_len = len(ddf.index)
-
- if row_len == 0:
- LOG.warning("Skipping, don't partial save an empty dd.DataFrame")
- return False
-
- try:
- assert schema.metadata
- partition = schema.metadata.get(PARTITION_ON)
- ddf.to_parquet(
- path=next_numbered_path,
- partition_on=partition,
- engine="pyarrow",
- overwrite=True,
- write_metadata_file=True,
- compression="brotli",
- )
- except (pa.ArrowInvalid, pa.ArrowIOError, OSError) as e:
- LOG.exception(e)
- self.delete_archive(next_numbered_path)
- return False
-
- if platform == "darwin":
- subprocess.call(["ln", "-sfn", target, partial_path])
- else:
- subprocess.call(["ln", "-sfnT", target, partial_path])
-
- return True
-
- def initial_load(self, overwrite: bool = False) -> bool:
-
- if overwrite is False:
- assert not self.has_archive(include_empty=True), "already archived"
-
- assert self.should_archive(), "not ready to archive!"
-
- df: pd.DataFrame | None = self.from_db()
-
- if df is None:
- self.set_empty()
- return False
-
- ddf = dd.from_pandas(df, npartitions=1)
- return self.to_archive(ddf=ddf, is_partial=False, overwrite=overwrite)
-
- def clear_corrupt_archive(self):
- if self.has_archive(include_empty=False) and not self.valid_archive(self.path):
- LOG.warning(f"invalid archive, deleting: {self.path}")
- self.delete_archive(self.path)
-
-
-class DFCollection(CollectionBase):
- data_type: DFCollectionType | None = Field(default=None)
-
- # --- Private ---
- pg_config: PostgresConfig | None = Field(default=None)
-
- def __repr__(self):
- res = self.signature() + "\n"
- if len(self.items) > 6:
- items = self.items[:3] + ["..."] + self.items[-3:]
- else:
- items = self.items
-
- for i in items:
- res += f" – {repr(i) if isinstance(i, DFCollectionItem) else i}\n"
-
- return res
-
- def signature(self):
- arr = [
- 1 if i.has_archive(include_empty=True) else 0
- for i in self.items
- if i.should_archive()
- ]
- repr_str = (
- f"items={len(self.items)}; start={self.start} @ {self.offset}; {int(sum(arr) / len(arr) * 100)}% "
- f"archived"
- )
- res = f"{self.__repr_name__()}({repr_str})"
- return res
-
- @field_validator("data_type")
- def check_data_type(cls, data_type: DFCollectionType | None, info: ValidationInfo):
- if data_type is None:
- raise ValueError("Must explicitly provide a data_type")
-
- if data_type not in DFCollectionTypeSchemas:
- raise ValueError("Must provide a supported data_type")
-
- return data_type
-
- # --- Properties ---
- @property
- def items(self) -> list[DFCollectionItem]:
- items = []
- for iv in self.interval_range:
- cm = DFCollectionItem(start=iv[0])
- cm._collection = self
- items.append(cm)
- return items
-
- @property
- def _schema(self) -> DataFrameSchema:
- return DFCollectionTypeSchemas[self.data_type]
-
- # --- Methods ---
-
- def initial_load(
- self,
- client: DaskClient | None = None,
- sync: bool = True,
- since: datetime | None = None,
- client_resources: dict[str, Any] | None = None,
- timeout: float | None = None,
- ) -> list[Future]:
- # This can be used to just build all local archive files
- # We typically want to go backwards first, so we can most quickly
- # populate the last 90 days for example
-
- client = client or self._client
-
- LOG.info(f"{self.data_type.value}.initial_load({since=}, {sync=})")
-
- items = self.items
- if since:
- items = self.get_items(since=since)
-
- if client is None:
- for item in reversed(items):
- if item.has_archive(include_empty=True):
- continue
- if not item.should_archive():
- continue
- item.initial_load()
- return []
-
- fs = []
- for item in items:
- if item.has_archive(include_empty=True):
- continue
- if not item.should_archive():
- continue
- f = dask.delayed(item.initial_load)()
- fs.append(f)
-
- if sync:
- fs = client.compute(fs, sync=False, priority=2, resources=client_resources)
- _ = as_completed(fs, timeout=timeout)
- return fs
-
- else:
- return client.compute(fs, sync=True, priority=2, resources=client_resources)
-
- def fetch_force_rr_latest(self, sources) -> list[FilePath]:
- LOG.info(
- f"{self.data_type.value}.fetch_force_rr_latest(sources={len(sources)})"
- )
-
- # We only want 'partial-able' items (those that can not yet be archived).
- rr_items = [
- i for i in self.items if not i.should_archive() and not i.is_empty()
- ]
- if rr_items:
- # If the ARCHIVE_AFTER time is > the collection offset (which it is always currently),
- # then there typically wouldn't be more than 1 un-archivable item.
- _start = rr_items[0].start
- _end = rr_items[-1].finish
- rr_duration = (_end - _start).total_seconds()
-
- # TODO: Do we want to be smarter about any rr selects max durations?
- # allowing 2x the length of the offset. If we have more than this not archived,
- # we want to run the archive first, not fetch from rr
- archive_after = self._schema.metadata[ARCHIVE_AFTER]
- allowed_rr_duration = (
- (pd.Timedelta(self.offset) * 2) + archive_after
- ).total_seconds()
- if rr_duration > allowed_rr_duration:
- raise ValueError(
- f"rr select duration exceeds {pd.Timedelta(allowed_rr_duration)}"
- )
-
- for rr_item in rr_items:
- if (
- rr_item.has_partial_archive()
- and self.data_type != DFCollectionType.LEDGER
- ):
- saved = rr_item.update_partial_archive()
- else:
- saved = rr_item.create_partial_archive()
- if saved:
- sources.append(rr_item.partial_path)
-
- return sources
-
- def force_rr_latest(
- self,
- client: DaskClient,
- client_resources: dict[str, Any] | None = None,
- sync: bool = True,
- ) -> list[Future]:
-
- # For forcing update of any partials asynchronously if desired
- LOG.info(f"{self.data_type.value}.force_rr_latest({client=})")
-
- rr_items = [
- i for i in self.items if not i.should_archive() and not i.is_empty()
- ]
- fs = []
- for rr_item in rr_items:
- if (
- rr_item.has_partial_archive()
- and self.data_type != DFCollectionType.LEDGER
- ):
- fs.append(dask.delayed(rr_item.update_partial_archive)())
- else:
- fs.append(dask.delayed(rr_item.create_partial_archive)())
- return client.compute(fs, sync=sync, priority=2, resources=client_resources)
diff --git a/generalresearch/incite/collections/base.py b/generalresearch/incite/collections/base.py
new file mode 100644
index 0000000..ecfa57a
--- /dev/null
+++ b/generalresearch/incite/collections/base.py
@@ -0,0 +1,693 @@
+from __future__ import annotations
+
+import os
+import subprocess
+import time
+from datetime import datetime
+from enum import StrEnum
+from sys import platform
+from typing import Any
+
+import dask
+import dask.dataframe as dd
+import pandas as pd
+import pyarrow as pa
+import pyarrow.parquet as pq
+from dask.distributed import Client as DaskClient
+from dask.distributed import Future
+from distributed import as_completed
+from more_itertools import chunked
+from pandera.pandas import DataFrameSchema
+from psycopg import Cursor
+from pydantic import Field, FilePath, ValidationInfo, field_validator
+from sentry_sdk import capture_exception
+
+from generalresearch.incite import LOG
+from generalresearch.incite.base import CollectionBase, CollectionItemBase
+from generalresearch.incite.schemas import (
+ ARCHIVE_AFTER,
+ ORDER_KEY,
+ PARTITION_ON,
+ empty_dataframe_from_schema,
+)
+from generalresearch.incite.schemas.thl_marketplaces import (
+ InnovateSurveyHistorySchema,
+ MorningSurveyTimeseriesSchema,
+ SagoSurveyHistorySchema,
+ SpectrumSurveyTimeseriesSchema,
+)
+from generalresearch.incite.schemas.thl_web import (
+ LedgerSchema,
+ THLIPInfoSchema,
+ THLSessionSchema,
+ THLTaskAdjustmentSchema,
+ THLUserSchema,
+ THLWallSchema,
+ TransactionMetadataColumns,
+ TxMetaSchema,
+ TxSchema,
+ UserHealthAuditLogSchema,
+ UserHealthIPHistorySchema,
+ UserHealthIPHistoryWSSchema,
+)
+from generalresearch.pg_helper import PostgresConfig
+
+DT_STR = "%Y-%m-%d %H:%M:%S"
+
+
+class DFCollectionType(StrEnum):
+ TEST = "test"
+
+ USER = "thl_user"
+ SESSION = "thl_session"
+ WALL = "thl_wall"
+ TASK_ADJUSTMENT = "thl_taskadjustment"
+ IP_INFO = "thl_ipinformation"
+
+ AUDIT_LOG = "userhealth_auditlog"
+ IP_HISTORY = "userhealth_iphistory"
+ IP_HISTORY_WS = "userhealth_iphistory_ws"
+
+ LEDGER = "ledger"
+
+ INNOVATE_SURVEY_HISTORY = "innovate_surveyhistory"
+ MORNING_SURVEY_TIMESERIES = "morning_surveytimeseries"
+ SAGO_SURVEY_HISTORY = "sago_surveyhistory"
+ SPECTRUM_SURVEY_TIMESERIES = "spectrum_surveytimeseries"
+
+
+DFCollectionTypeSchemas = {
+ DFCollectionType.USER: THLUserSchema,
+ DFCollectionType.WALL: THLWallSchema,
+ DFCollectionType.SESSION: THLSessionSchema,
+ DFCollectionType.IP_INFO: THLIPInfoSchema,
+ DFCollectionType.TASK_ADJUSTMENT: THLTaskAdjustmentSchema,
+ DFCollectionType.IP_HISTORY: UserHealthIPHistorySchema,
+ DFCollectionType.IP_HISTORY_WS: UserHealthIPHistoryWSSchema,
+ DFCollectionType.AUDIT_LOG: UserHealthAuditLogSchema,
+ DFCollectionType.LEDGER: LedgerSchema,
+ DFCollectionType.INNOVATE_SURVEY_HISTORY: InnovateSurveyHistorySchema,
+ DFCollectionType.MORNING_SURVEY_TIMESERIES: MorningSurveyTimeseriesSchema,
+ DFCollectionType.SAGO_SURVEY_HISTORY: SagoSurveyHistorySchema,
+ DFCollectionType.SPECTRUM_SURVEY_TIMESERIES: SpectrumSurveyTimeseriesSchema,
+}
+
+
+class DFCollectionItem(CollectionItemBase):
+
+ # --- Properties ---
+ @property
+ def filename(self) -> str:
+ return (
+ f"{self._collection.data_type.name.lower()}-{self._collection.offset}"
+ f"-{self.start.strftime('%Y-%m-%d-%H-%M-%S')}.parquet"
+ )
+
+ # --- Methods ---
+
+ def has_postgres(self) -> bool:
+ if self._collection.pg_config is None:
+ return False
+
+ connected = True
+ try:
+ self._collection.pg_config.execute_sql_query("""SELECT 1;""")
+ except AssertionError:
+ connected = False
+
+ return connected
+
+ def has_db(self) -> bool:
+ return self.has_mysql() or self.has_postgres()
+
+ def update_partial_archive(self) -> bool:
+ if not self.valid_archive(self.partial_path, sample=1000):
+ LOG.error(f"invalid partial archive: {self.partial_path}")
+ return self.create_partial_archive()
+ df = pq.ParquetDataset(self.partial_path).read().to_pandas()
+
+ order_key = self._collection._schema.metadata[ORDER_KEY]
+ archive_after = self._collection._schema.metadata[ARCHIVE_AFTER]
+
+ partial_max = df[order_key].max().to_pydatetime()
+
+ since = partial_max - archive_after
+ since = max([since, self.start]) # don't allow to query before the item's start
+ df = df[df[order_key] < since].copy()
+
+ _df = self.from_db(since=since)
+
+ if _df is not None:
+ df = pd.concat([df, _df])
+ self.to_archive(ddf=dd.from_pandas(df, npartitions=1), is_partial=True)
+ else:
+ # The update to the partial returned no rows, but the partial
+ # still exists, so we'll continue with whatever was calling this.
+ # We don't need to re-write the partial or really do anything.
+ pass
+ return True
+
+ def create_partial_archive(self) -> bool:
+ _df = self.from_db()
+ if _df is None:
+ # Returned no rows, but the period is not closed, so we
+ # don't want to mark as empty. Do nothing.
+ return False
+ return self.to_archive(ddf=dd.from_pandas(_df, npartitions=1), is_partial=True)
+
+ # --- ORM / Data handlers---
+ def to_dict(self) -> dict[str, Any]:
+ return self._to_dict()
+
+ def from_db(self, since: datetime | None = None) -> pd.DataFrame | None:
+ if self._collection.data_type == DFCollectionType.LEDGER:
+ assert since is None, "Shouldn't pass since for Ledger item"
+ assert self._collection.pg_config is not None
+ return self.from_postgres_ledger()
+ else:
+ return self.from_db_standard(since=since)
+
+ def from_db_standard(self, since: datetime | None = None) -> pd.DataFrame | None:
+ assert (
+ self._collection.data_type != DFCollectionType.LEDGER
+ ), "Can't call from_postgres_standard for Ledger DFCollectionItem"
+
+ start, finish = self.start, self.finish
+ LOG.debug(
+ f"{self._collection.data_type.value}.from_postgres("
+ f"start={start.strftime(DT_STR)}, "
+ f"finish={finish.strftime(DT_STR)})"
+ )
+ coll = self._collection
+ schema = coll._schema
+ pg_config = coll.pg_config
+ assert pg_config, "Must provide PostgresConfig"
+
+ start = since or start
+ order_key = schema.metadata[ORDER_KEY]
+ cols = list(schema.columns.keys()) + [schema.index.name]
+ cols_str = ", ".join(cols)
+
+ try:
+ res = pg_config.execute_sql_query(
+ query=f"""
+ SELECT {cols_str}
+ FROM {coll.data_type.value}
+ WHERE {order_key} >= %s AND {order_key} < %s;
+ """,
+ params=[start, finish],
+ )
+ except AssertionError as e:
+ capture_exception(error=e)
+ LOG.error(f"_from_postgres Exception: {e}")
+ return None
+
+ if not res:
+ LOG.warning("_from_postgres query returned nothing")
+ # Return an empty df.DataFrame with the correct columns
+ return empty_dataframe_from_schema(coll._schema)
+
+ df = pd.DataFrame.from_records(res).set_index(coll._schema.index.name)
+ df = self.validate_df(df=df)
+
+ if df is None:
+ LOG.warning("_from_postgres query results failed validation")
+ # Schema validation can fail...
+ return None
+
+ return df
+
+ def from_postgres_ledger(self) -> pd.DataFrame | None:
+ assert (
+ self._collection.data_type == DFCollectionType.LEDGER
+ ), "Can only call from_postgres_ledger on Ledger DFCollectionItem"
+
+ start, finish = self.start, self.finish
+ LOG.info(
+ f"{self._collection.data_type.value}.from_postgres_ledger("
+ f"start={start.strftime(DT_STR)}, "
+ f"finish={finish.strftime(DT_STR)})"
+ )
+
+ coll = self._collection
+ assert coll.pg_config, "Must provide PostgresConfig"
+ pg_config: PostgresConfig = coll.pg_config
+
+ limit = 20000
+ offset = 0
+ res = []
+ while True:
+ LOG.info(
+ f"{self._collection.data_type.value}.from_postgres_ledger({limit=}, {offset=})"
+ )
+ chunk = pg_config.execute_sql_query(
+ query=f"""
+ SELECT lt.id AS tx_id, lt.created, lt.ext_description, lt.tag,
+ le.id AS entry_id, le.direction, le.amount, le.account_id,
+ la.display_name, la.qualified_name, la.account_type,
+ la.normal_balance, la.reference_type, la.reference_uuid,
+ la.currency
+ FROM ledger_transaction AS lt
+ LEFT JOIN ledger_entry AS le
+ ON lt.id = le.transaction_id
+ LEFT JOIN ledger_account AS la
+ ON la.uuid = le.account_id
+ WHERE lt.created >= %s AND lt.created < %s
+ AND le.id IS NOT NULL
+ ORDER BY lt.created
+ LIMIT {limit} OFFSET {offset};
+ """,
+ params=[start, finish],
+ )
+ res.extend(chunk)
+ if not chunk:
+ break
+ offset += limit
+
+ if len(res) == 0:
+ return None
+
+ # Note (AND le.id IS NOT NULL): It is possible we have transactions with
+ # no ledger entries. This is because the transaction creation failed
+ # for some reason. The ledger is not unbalanced, it is just an orphan
+ # transaction. Just skip those here.
+
+ tx_df = TxSchema.validate(
+ check_obj=pd.DataFrame.from_records(res).set_index("entry_id"),
+ lazy=True,
+ )
+
+ tx_ids = list(tx_df["tx_id"].unique())
+ metadata_res = []
+ # "MySQL server has gone away" if this is too big
+ conn = pg_config.make_connection()
+ c: Cursor = conn.cursor()
+ for chunk in chunked(tx_ids, n=5_000):
+ c.execute(
+ query="""
+ SELECT ltm.transaction_id AS tx_id,
+ ltm.id AS tx_metadata_id,
+ ltm.key, ltm.value
+ FROM ledger_transactionmetadata AS ltm
+ WHERE ltm.transaction_id = ANY(%s);
+ """,
+ params=[chunk],
+ )
+ metadata_res += c.fetchall()
+
+ conn.close()
+
+ tx_meta = (
+ pd.DataFrame(
+ TxMetaSchema.validate(
+ check_obj=pd.DataFrame.from_records(metadata_res).set_index(
+ ["tx_id", "tx_metadata_id"]
+ ),
+ lazy=True,
+ ).pivot(columns="key", values="value"),
+ # This makes sure we expand to have all the possible columns
+ columns=[e.value for e in TransactionMetadataColumns],
+ )
+ .groupby("tx_id")
+ .first()
+ )
+
+ df = tx_df.merge(tx_meta, how="left", left_on="tx_id", right_index=True)
+ df = self.validate_df(df=df)
+
+ if df is None:
+ # Schema validation can fail...
+ return None
+
+ return df
+
+ def to_archive(
+ self,
+ ddf: dd.DataFrame,
+ is_partial: bool = False,
+ overwrite: bool = False,
+ ) -> bool:
+ """
+ :returns: bool (saved_successful)
+ """
+ assert isinstance(ddf, dd.DataFrame), "must pass dask df"
+
+ client: DaskClient | None = self._collection._client
+ # client = None
+
+ if client:
+ row_len = client.compute(collections=ddf.shape[0], sync=True)
+ else:
+ row_len = len(ddf.index)
+ is_empty = row_len == 0
+
+ if is_partial:
+ return self.to_archive_numbered_partial(ddf=ddf)
+ else:
+ return self._to_archive(
+ ddf=ddf,
+ is_empty=is_empty,
+ overwrite=overwrite,
+ )
+
+ def _to_archive(
+ self,
+ ddf: dd.DataFrame | None,
+ is_empty: bool,
+ overwrite: bool = False,
+ ) -> bool:
+ """
+ For archiving an item. Will write an empty file if ddf is empty. This
+ is NOT for writing partials.
+
+ :returns: bool (saved_successful)
+ """
+
+ if ddf is None:
+ return False
+
+ should_archive = self.should_archive()
+ if not should_archive:
+ LOG.warning(f"Cannot create archive for such new data: {self.path}")
+ return False
+
+ if overwrite is False:
+ has_archive = self.has_archive(include_empty=True)
+ if has_archive:
+ LOG.warning(f"archive already exists: {self.path}")
+ return False
+
+ if is_empty:
+ # Create an .empty only if the Item is "archiveable" (which we checked above)
+ self.set_empty()
+ return True
+
+ # Incase the file saving is interrupted, or otherwise fails
+ # save it to a tmp file first, then rename once we can confirm
+ # that it successfully loads
+ tmp_path = self.tmp_path()
+ try:
+ schema = self._collection._schema
+ assert schema
+ assert schema.metadata
+ partition = schema.metadata.get(PARTITION_ON)
+
+ ddf.to_parquet(
+ path=tmp_path,
+ partition_on=partition,
+ engine="pyarrow",
+ overwrite=True,
+ write_metadata_file=True,
+ compression="brotli",
+ )
+
+ except (pa.ArrowInvalid, pa.ArrowIOError, OSError) as e:
+ LOG.exception(e)
+ self.delete_archive(tmp_path)
+ return False
+
+ # It was saved, but the file seems to be corrupt
+ if not self.valid_archive(tmp_path):
+ LOG.error(f"not valid archive: {tmp_path}")
+ self.delete_archive(tmp_path)
+ # File did not save correctly so return it as saved=False
+ return False
+
+ # To debug, just set this key to auto expire in 5 seconds
+ # RC.set(name=f"_to_archive:{self.path.as_posix()}", value=1, ex=15)
+ # with RC.lock(f"_to_archive:{self.path.as_posix()}:lock", timeout=15):
+
+ if os.path.isfile(tmp_path):
+ # If the file was saved okay, seems okay, rename it
+ os.replace(tmp_path, self.path)
+ os.remove(tmp_path)
+
+ if os.path.isdir(tmp_path):
+ if os.path.exists(self.path.as_posix()):
+ if overwrite:
+ subprocess.call(["rm", "-r", self.path.as_posix()])
+ time.sleep(1)
+ else:
+ LOG.error(f"already exists: {self.path.as_posix()}")
+ return False
+
+ if platform == "darwin":
+ subprocess.call(["mv", tmp_path.as_posix(), self.path.as_posix()])
+ else:
+ # -T will (should) cause the mv to fail if path wasn't successfully deleted
+ subprocess.call(["mv", "-T", tmp_path.as_posix(), self.path.as_posix()])
+ return True
+
+ def to_archive_numbered_partial(self, ddf: dd.DataFrame | None = None) -> bool:
+ """
+ For partial files/dirs only. Writes the .partial file with a number
+ at the end (.partial.####) and then creates a symlink
+ from .partial -> .partial.####
+
+ :returns: bool (saved_successful)
+ """
+ if ddf is None:
+ return False
+
+ collection = self._collection
+ schema = collection._schema
+ assert schema
+ client: DaskClient | None = collection._client
+
+ next_numbered_path = self.next_numbered_path(self.partial_path)
+ partial_path = self.partial_path
+ # finish = self.finish
+
+ # Make sure these are in the same dir. b/c the symlink has to be
+ # relative, not an absolute path
+ assert (
+ partial_path.parent == next_numbered_path.parent
+ ), "Can't have numbered_path in a different directory"
+ target = (
+ next_numbered_path.name
+ ) # this is the symlink's target. it is a relative path (only the name)
+
+ should_archive = self.should_archive()
+ assert should_archive is False, "Don't write partial if the item is archiveable"
+
+ if client:
+ row_len = client.compute(collections=ddf.shape[0], sync=True)
+ else:
+ row_len = len(ddf.index)
+
+ if row_len == 0:
+ LOG.warning("Skipping, don't partial save an empty dd.DataFrame")
+ return False
+
+ try:
+ assert schema.metadata
+ partition = schema.metadata.get(PARTITION_ON)
+ ddf.to_parquet(
+ path=next_numbered_path,
+ partition_on=partition,
+ engine="pyarrow",
+ overwrite=True,
+ write_metadata_file=True,
+ compression="brotli",
+ )
+ except (pa.ArrowInvalid, pa.ArrowIOError, OSError) as e:
+ LOG.exception(e)
+ self.delete_archive(next_numbered_path)
+ return False
+
+ if platform == "darwin":
+ subprocess.call(["ln", "-sfn", target, partial_path])
+ else:
+ subprocess.call(["ln", "-sfnT", target, partial_path])
+
+ return True
+
+ def initial_load(self, overwrite: bool = False) -> bool:
+
+ if overwrite is False:
+ assert not self.has_archive(include_empty=True), "already archived"
+
+ assert self.should_archive(), "not ready to archive!"
+
+ df: pd.DataFrame | None = self.from_db()
+
+ if df is None:
+ self.set_empty()
+ return False
+
+ ddf = dd.from_pandas(df, npartitions=1)
+ return self.to_archive(ddf=ddf, is_partial=False, overwrite=overwrite)
+
+ def clear_corrupt_archive(self):
+ if self.has_archive(include_empty=False) and not self.valid_archive(self.path):
+ LOG.warning(f"invalid archive, deleting: {self.path}")
+ self.delete_archive(self.path)
+
+
+class DFCollection(CollectionBase):
+ data_type: DFCollectionType | None = Field(default=None)
+
+ # --- Private ---
+ pg_config: PostgresConfig | None = Field(default=None)
+
+ def __repr__(self):
+ res = self.signature() + "\n"
+ if len(self.items) > 6:
+ items = self.items[:3] + ["..."] + self.items[-3:]
+ else:
+ items = self.items
+
+ for i in items:
+ res += f" – {repr(i) if isinstance(i, DFCollectionItem) else i}\n"
+
+ return res
+
+ def signature(self):
+ arr = [
+ 1 if i.has_archive(include_empty=True) else 0
+ for i in self.items
+ if i.should_archive()
+ ]
+ repr_str = (
+ f"items={len(self.items)}; start={self.start} @ {self.offset}; {int(sum(arr) / len(arr) * 100)}% "
+ f"archived"
+ )
+ res = f"{self.__repr_name__()}({repr_str})"
+ return res
+
+ @field_validator("data_type")
+ def check_data_type(cls, data_type: DFCollectionType | None, info: ValidationInfo):
+ if data_type is None:
+ raise ValueError("Must explicitly provide a data_type")
+
+ if data_type not in DFCollectionTypeSchemas:
+ raise ValueError("Must provide a supported data_type")
+
+ return data_type
+
+ # --- Properties ---
+ @property
+ def items(self) -> list[DFCollectionItem]:
+ items = []
+ for iv in self.interval_range:
+ cm = DFCollectionItem(start=iv[0])
+ cm._collection = self
+ items.append(cm)
+ return items
+
+ @property
+ def _schema(self) -> DataFrameSchema:
+ return DFCollectionTypeSchemas[self.data_type]
+
+ # --- Methods ---
+
+ def initial_load(
+ self,
+ client: DaskClient | None = None,
+ sync: bool = True,
+ since: datetime | None = None,
+ client_resources: dict[str, Any] | None = None,
+ timeout: float | None = None,
+ ) -> list[Future]:
+ # This can be used to just build all local archive files
+ # We typically want to go backwards first, so we can most quickly
+ # populate the last 90 days for example
+
+ client = client or self._client
+
+ LOG.info(f"{self.data_type.value}.initial_load({since=}, {sync=})")
+
+ items = self.items
+ if since:
+ items = self.get_items(since=since)
+
+ if client is None:
+ for item in reversed(items):
+ if item.has_archive(include_empty=True):
+ continue
+ if not item.should_archive():
+ continue
+ item.initial_load()
+ return []
+
+ fs = []
+ for item in items:
+ if item.has_archive(include_empty=True):
+ continue
+ if not item.should_archive():
+ continue
+ f = dask.delayed(item.initial_load)()
+ fs.append(f)
+
+ if sync:
+ fs = client.compute(fs, sync=False, priority=2, resources=client_resources)
+ _ = as_completed(fs, timeout=timeout)
+ return fs
+
+ else:
+ return client.compute(fs, sync=True, priority=2, resources=client_resources)
+
+ def fetch_force_rr_latest(self, sources) -> list[FilePath]:
+ LOG.info(
+ f"{self.data_type.value}.fetch_force_rr_latest(sources={len(sources)})"
+ )
+
+ # We only want 'partial-able' items (those that can not yet be archived).
+ rr_items = [
+ i for i in self.items if not i.should_archive() and not i.is_empty()
+ ]
+ if rr_items:
+ # If the ARCHIVE_AFTER time is > the collection offset (which it is always currently),
+ # then there typically wouldn't be more than 1 un-archivable item.
+ _start = rr_items[0].start
+ _end = rr_items[-1].finish
+ rr_duration = (_end - _start).total_seconds()
+
+ # TODO: Do we want to be smarter about any rr selects max durations?
+ # allowing 2x the length of the offset. If we have more than this not archived,
+ # we want to run the archive first, not fetch from rr
+ archive_after = self._schema.metadata[ARCHIVE_AFTER]
+ allowed_rr_duration = (
+ (pd.Timedelta(self.offset) * 2) + archive_after
+ ).total_seconds()
+ if rr_duration > allowed_rr_duration:
+ raise ValueError(
+ f"rr select duration exceeds {pd.Timedelta(allowed_rr_duration)}"
+ )
+
+ for rr_item in rr_items:
+ if (
+ rr_item.has_partial_archive()
+ and self.data_type != DFCollectionType.LEDGER
+ ):
+ saved = rr_item.update_partial_archive()
+ else:
+ saved = rr_item.create_partial_archive()
+ if saved:
+ sources.append(rr_item.partial_path)
+
+ return sources
+
+ def force_rr_latest(
+ self,
+ client: DaskClient,
+ client_resources: dict[str, Any] | None = None,
+ sync: bool = True,
+ ) -> list[Future]:
+
+ # For forcing update of any partials asynchronously if desired
+ LOG.info(f"{self.data_type.value}.force_rr_latest({client=})")
+
+ rr_items = [
+ i for i in self.items if not i.should_archive() and not i.is_empty()
+ ]
+ fs = []
+ for rr_item in rr_items:
+ if (
+ rr_item.has_partial_archive()
+ and self.data_type != DFCollectionType.LEDGER
+ ):
+ fs.append(dask.delayed(rr_item.update_partial_archive)())
+ else:
+ fs.append(dask.delayed(rr_item.create_partial_archive)())
+ return client.compute(fs, sync=sync, priority=2, resources=client_resources)
diff --git a/generalresearch/incite/collections/thl_marketplaces.py b/generalresearch/incite/collections/thl_marketplaces.py
index fe2b01f..246fe87 100644
--- a/generalresearch/incite/collections/thl_marketplaces.py
+++ b/generalresearch/incite/collections/thl_marketplaces.py
@@ -1,6 +1,6 @@
from typing import Literal
-from generalresearch.incite.collections import DFCollection, DFCollectionType
+from generalresearch.incite.collections.base import DFCollection, DFCollectionType
from generalresearch.incite.schemas.thl_marketplaces import (
InnovateSurveyHistorySchema,
MorningSurveyTimeseriesSchema,
diff --git a/generalresearch/incite/defaults.py b/generalresearch/incite/defaults.py
index d4025fc..773555c 100644
--- a/generalresearch/incite/defaults.py
+++ b/generalresearch/incite/defaults.py
@@ -15,7 +15,7 @@ from generalresearch.incite.collections.thl_web import (
UserDFCollection,
WallDFCollection,
)
-from generalresearch.incite.mergers import MergeType
+from generalresearch.incite.mergers.base import MergeType
from generalresearch.incite.mergers.foundations.enriched_session import (
EnrichedSessionMerge,
)
@@ -47,9 +47,7 @@ def session_df_collection(
)
-def wall_df_collection(
- ds: GRLDatasets, pg_config: PostgresConfig
-) -> WallDFCollection:
+def wall_df_collection(ds: GRLDatasets, pg_config: PostgresConfig) -> WallDFCollection:
return WallDFCollection(
offset="49h",
pg_config=pg_config,
@@ -58,9 +56,7 @@ def wall_df_collection(
)
-def user_df_collection(
- ds: GRLDatasets, pg_config: PostgresConfig
-) -> UserDFCollection:
+def user_df_collection(ds: GRLDatasets, pg_config: PostgresConfig) -> UserDFCollection:
return UserDFCollection(
offset="73h",
pg_config=pg_config,
diff --git a/generalresearch/incite/mergers/__init__.py b/generalresearch/incite/mergers/__init__.py
index 810bedc..e69de29 100644
--- a/generalresearch/incite/mergers/__init__.py
+++ b/generalresearch/incite/mergers/__init__.py
@@ -1,301 +0,0 @@
-import logging
-import os.path
-import subprocess
-from datetime import UTC, datetime
-from enum import StrEnum
-from sys import platform
-from typing import Self
-
-import dask.dataframe as dd
-import pandas as pd
-from dask.distributed import Client
-from pandera.pandas import DataFrameSchema
-from pydantic import Field, ValidationInfo, field_validator, model_validator
-
-from generalresearch.incite.base import CollectionBase, CollectionItemBase
-from generalresearch.incite.schemas import PARTITION_ON
-from generalresearch.incite.schemas.mergers.foundations.enriched_session import (
- EnrichedSessionSchema,
-)
-from generalresearch.incite.schemas.mergers.foundations.enriched_task_adjust import (
- EnrichedTaskAdjustSchema,
-)
-from generalresearch.incite.schemas.mergers.foundations.enriched_wall import (
- EnrichedWallSchema,
-)
-from generalresearch.incite.schemas.mergers.foundations.user_id_product import (
- UserIdProductSchema,
-)
-from generalresearch.incite.schemas.mergers.pop_ledger import (
- PopLedgerSchema,
-)
-from generalresearch.incite.schemas.mergers.ym_survey_wall import (
- YMSurveyWallSchema,
-)
-from generalresearch.incite.schemas.mergers.ym_wall_summary import (
- YMWallSummarySchema,
-)
-from generalresearch.models.custom_types import AwareDatetimeISO
-
-LOG = logging.getLogger("incite")
-
-
-class MergeType(StrEnum):
- TEST = "test"
- YM_SURVEY_WALL = "ym_survey_wall"
- YM_WALL_SUMMARY = "ym_wall_summary"
-
- POP_LEDGER = "pop_ledger"
-
- # --- Foundations ---
- USER_ID_PRODUCT = "user_id_product"
- ENRICHED_WALL = "enriched_wall"
- ENRICHED_SESSION = "enriched_session"
- ENRICHED_TASK_ADJUST = "enriched_task_adjust"
-
-
-MergeTypeSchemas = {
- MergeType.YM_SURVEY_WALL: YMSurveyWallSchema,
- MergeType.YM_WALL_SUMMARY: YMWallSummarySchema,
- MergeType.POP_LEDGER: PopLedgerSchema,
- # --- Foundations ---
- MergeType.USER_ID_PRODUCT: UserIdProductSchema,
- MergeType.ENRICHED_WALL: EnrichedWallSchema,
- MergeType.ENRICHED_SESSION: EnrichedSessionSchema,
- MergeType.ENRICHED_TASK_ADJUST: EnrichedTaskAdjustSchema,
-}
-
-
-class MergeCollectionItem(CollectionItemBase):
-
- # --- Properties ---
-
- @property
- def finish(self) -> datetime:
- # A MergeCollection can have offset = None
- if self._collection.offset:
- return (
- pd.Timestamp(self.start) + pd.Timedelta(self._collection.offset)
- ).to_pydatetime()
- else:
- return datetime.now(tz=UTC).replace(microsecond=0)
-
- @property
- def filename(self) -> str:
- grouped_key = self._collection.grouped_key
- offset = self._collection.offset
- start = self.start.strftime("%Y-%m-%d-%H-%M-%S")
- f = [self._collection.merge_type.name.lower()]
- if offset:
- f.append(offset)
- if grouped_key:
- f.append(grouped_key)
- if self._collection.start is not None:
- # This is a collection that is "looking back" 'offset' time (1 item).
- f.append(start)
- s = "-".join(f)
- s += ".parquet"
- return s
-
- # --- ORM / Data handlers---
- def to_dict(self, *args, **kwargs) -> dict:
- res = self._to_dict()
- res["group_by"] = self._collection.group_by
- return res
-
- def to_archive(
- self,
- client: Client,
- ddf: dd.DataFrame,
- is_partial: bool = False,
- ) -> bool:
- assert is_partial is False, "use to_archive_symlink"
- return self._to_archive(client=client, ddf=ddf, client_resources=None)
-
- def _to_archive(
- self, client: Client, ddf: dd.DataFrame | None, client_resources=None
- ) -> bool:
- """
- For archiving an item. Will write an empty file if ddf is empty.
- This is NOT for writing partials.
-
- :returns: bool (saved_successful)
- """
- if ddf is None:
- return False
-
- row_len: int = client.compute(collections=ddf.shape[0], sync=True)
- assert row_len
- assert row_len > 0, "empty ddf"
-
- tmp_path = self.tmp_path()
- schema = self._collection._schema
- assert schema.metadata
-
- partition = schema.metadata.get(PARTITION_ON)
- f = ddf.to_parquet(
- compute=False,
- path=tmp_path,
- partition_on=partition,
- engine="pyarrow",
- overwrite=True,
- write_metadata_file=True,
- compression="brotli",
- )
- client.compute(f, sync=True, priority=2, resources=client_resources)
- assert not os.path.exists(
- self.path.as_posix()
- ), f"already exits!: {self.path.as_posix()}"
-
- if platform == "darwin":
- subprocess.call(["mv", tmp_path.as_posix(), self.path.as_posix()])
- else:
- # -T will (should) cause the mv to fail if `path` wasn't successfully deleted
- subprocess.call(["mv", "-T", tmp_path.as_posix(), self.path.as_posix()])
- return True
-
- def to_archive_symlink(
- self,
- client: Client,
- ddf: dd.DataFrame,
- is_partial: bool = False,
- client_resources=None,
- validate_after=True,
- ) -> bool:
- """
- This differs from to_archive():
- 1) to_parquet is run in this process. If the df is already
- computed, there is no point in sending it to another worker
- to write.
-
- 2) symlink to next_numbered_path is created whether or not
- is_partial (to_archive only does this on partials)
-
- 3) we do not validate the written file. seems not useful to do
- this, as the file will probably get overwritten on the next
- loop anyway
- """
- path = self.partial_path if is_partial else self.path
- next_numbered_path = self.next_numbered_path(path)
- collection = self._collection
- LOG.warning(f"{collection.merge_type.value}.to_archive_symlink()")
-
- assert isinstance(ddf, dd.DataFrame), "must pass a dask df"
-
- # We should validate before or after!!!
- # _validate_df(self.compute(ddf), coll._schema)
- target = (
- next_numbered_path.name
- ) # this is the symlink's target. it is a relative path (only the name)
-
- schema = self._collection._schema
- partition = schema.metadata.get(PARTITION_ON, None)
- f = ddf.to_parquet(
- compute=False,
- path=next_numbered_path.as_posix(),
- partition_on=partition,
- engine="pyarrow",
- overwrite=True,
- write_metadata_file=True,
- compression="brotli",
- )
- client.compute(f, sync=True, priority=2, resources=client_resources)
-
- if os.path.exists(path.as_posix()) and not os.path.islink(path.as_posix()):
- # This will fail when going from the old way to using symlinks,
- # if self.path already exists and is a directory.
- raise ValueError(
- f"first time we run this, make sure the path doesnt exist: {path.as_posix()}"
- )
-
- if platform == "darwin":
- subprocess.call(["ln", "-sfn", target, path.as_posix()])
- else:
- subprocess.call(["ln", "-sfnT", target, path.as_posix()])
-
- if validate_after and not self.valid_archive(self.path):
- LOG.error(f"{collection.merge_type.value} failed validation: {self.path}")
- self.delete_archive(self.path)
- return False
- return True
-
- # todo: unclear what the common interface should be here ... ?
- def fetch(self, *args, **kwargs) -> pd.DataFrame | dd.DataFrame:
- raise NotImplementedError("implement in subclass")
-
- def build(self, *args, **kwargs) -> pd.DataFrame | dd.DataFrame:
- raise NotImplementedError("implement in subclass")
-
-
-class MergeCollection(CollectionBase):
- """Mergers take instances of DFCollections, and/or other Mergers"""
-
- # In a merge, we can set offset = None which indicates that there is only 1
- # period/item where the range is 'start' until now.
- offset: str | None = Field(default="72h")
- # In a merge, we can set start = None which indicates that there is only 1
- # period/item where the range is (now - offset) until now.
- start: AwareDatetimeISO | None = Field(
- default=None,
- description="This is the starting point in which data will"
- " be retrieved in chunks from.",
- frozen=True,
- )
-
- merge_type: MergeType | None = Field(default=None)
- group_by: str | None = Field(default=None)
- grouped_key: str | None = Field(default=None)
- collection_item_class: type[MergeCollectionItem] = MergeCollectionItem
-
- @model_validator(mode="after")
- def check_start_and_offset_nullable(self) -> Self:
- if self.offset is None and self.start is None:
- raise AssertionError("cannot set both start and offset to None")
- return self
-
- @field_validator("merge_type")
- def check_merge_type(cls, merge_type: MergeType | None, info: ValidationInfo):
- if merge_type is None:
- raise ValueError("Must explicitly provide a merge_type")
-
- if merge_type not in MergeTypeSchemas:
- raise ValueError("Must provide a supported merge_type")
-
- return merge_type
-
- # --- Properties ---
- @property
- def interval_start(self) -> datetime | None:
- # if self.start is None and self.offset is set, the inferred start is (now - offset)
- if self.start is None:
- return datetime.now(tz=UTC).replace(microsecond=0) - pd.Timedelta(
- self.offset
- )
- return self.start
-
- @property
- def items(self) -> list[MergeCollectionItem]:
- items = []
- for iv in self.interval_range:
- cm = self.collection_item_class(start=iv[0])
- cm._collection = self
- items.append(cm)
- return items
-
- @property
- def _schema(self) -> DataFrameSchema:
- return MergeTypeSchemas[self.merge_type]
-
- def signature(self) -> str:
- arr = [
- 1 if i.has_archive(include_empty=True) else 0
- for i in self.items
- if i.should_archive()
- ]
- repr_str = (
- f"path={self.archive_path.as_posix()}; "
- f"items={len(self.items)}; start={self.start} @ {self.offset}; {int(sum(arr) / len(arr) * 100)}% "
- f"archived"
- )
- res = f"{self.__repr_name__()}({repr_str})"
- return res
diff --git a/generalresearch/incite/mergers/base.py b/generalresearch/incite/mergers/base.py
new file mode 100644
index 0000000..810bedc
--- /dev/null
+++ b/generalresearch/incite/mergers/base.py
@@ -0,0 +1,301 @@
+import logging
+import os.path
+import subprocess
+from datetime import UTC, datetime
+from enum import StrEnum
+from sys import platform
+from typing import Self
+
+import dask.dataframe as dd
+import pandas as pd
+from dask.distributed import Client
+from pandera.pandas import DataFrameSchema
+from pydantic import Field, ValidationInfo, field_validator, model_validator
+
+from generalresearch.incite.base import CollectionBase, CollectionItemBase
+from generalresearch.incite.schemas import PARTITION_ON
+from generalresearch.incite.schemas.mergers.foundations.enriched_session import (
+ EnrichedSessionSchema,
+)
+from generalresearch.incite.schemas.mergers.foundations.enriched_task_adjust import (
+ EnrichedTaskAdjustSchema,
+)
+from generalresearch.incite.schemas.mergers.foundations.enriched_wall import (
+ EnrichedWallSchema,
+)
+from generalresearch.incite.schemas.mergers.foundations.user_id_product import (
+ UserIdProductSchema,
+)
+from generalresearch.incite.schemas.mergers.pop_ledger import (
+ PopLedgerSchema,
+)
+from generalresearch.incite.schemas.mergers.ym_survey_wall import (
+ YMSurveyWallSchema,
+)
+from generalresearch.incite.schemas.mergers.ym_wall_summary import (
+ YMWallSummarySchema,
+)
+from generalresearch.models.custom_types import AwareDatetimeISO
+
+LOG = logging.getLogger("incite")
+
+
+class MergeType(StrEnum):
+ TEST = "test"
+ YM_SURVEY_WALL = "ym_survey_wall"
+ YM_WALL_SUMMARY = "ym_wall_summary"
+
+ POP_LEDGER = "pop_ledger"
+
+ # --- Foundations ---
+ USER_ID_PRODUCT = "user_id_product"
+ ENRICHED_WALL = "enriched_wall"
+ ENRICHED_SESSION = "enriched_session"
+ ENRICHED_TASK_ADJUST = "enriched_task_adjust"
+
+
+MergeTypeSchemas = {
+ MergeType.YM_SURVEY_WALL: YMSurveyWallSchema,
+ MergeType.YM_WALL_SUMMARY: YMWallSummarySchema,
+ MergeType.POP_LEDGER: PopLedgerSchema,
+ # --- Foundations ---
+ MergeType.USER_ID_PRODUCT: UserIdProductSchema,
+ MergeType.ENRICHED_WALL: EnrichedWallSchema,
+ MergeType.ENRICHED_SESSION: EnrichedSessionSchema,
+ MergeType.ENRICHED_TASK_ADJUST: EnrichedTaskAdjustSchema,
+}
+
+
+class MergeCollectionItem(CollectionItemBase):
+
+ # --- Properties ---
+
+ @property
+ def finish(self) -> datetime:
+ # A MergeCollection can have offset = None
+ if self._collection.offset:
+ return (
+ pd.Timestamp(self.start) + pd.Timedelta(self._collection.offset)
+ ).to_pydatetime()
+ else:
+ return datetime.now(tz=UTC).replace(microsecond=0)
+
+ @property
+ def filename(self) -> str:
+ grouped_key = self._collection.grouped_key
+ offset = self._collection.offset
+ start = self.start.strftime("%Y-%m-%d-%H-%M-%S")
+ f = [self._collection.merge_type.name.lower()]
+ if offset:
+ f.append(offset)
+ if grouped_key:
+ f.append(grouped_key)
+ if self._collection.start is not None:
+ # This is a collection that is "looking back" 'offset' time (1 item).
+ f.append(start)
+ s = "-".join(f)
+ s += ".parquet"
+ return s
+
+ # --- ORM / Data handlers---
+ def to_dict(self, *args, **kwargs) -> dict:
+ res = self._to_dict()
+ res["group_by"] = self._collection.group_by
+ return res
+
+ def to_archive(
+ self,
+ client: Client,
+ ddf: dd.DataFrame,
+ is_partial: bool = False,
+ ) -> bool:
+ assert is_partial is False, "use to_archive_symlink"
+ return self._to_archive(client=client, ddf=ddf, client_resources=None)
+
+ def _to_archive(
+ self, client: Client, ddf: dd.DataFrame | None, client_resources=None
+ ) -> bool:
+ """
+ For archiving an item. Will write an empty file if ddf is empty.
+ This is NOT for writing partials.
+
+ :returns: bool (saved_successful)
+ """
+ if ddf is None:
+ return False
+
+ row_len: int = client.compute(collections=ddf.shape[0], sync=True)
+ assert row_len
+ assert row_len > 0, "empty ddf"
+
+ tmp_path = self.tmp_path()
+ schema = self._collection._schema
+ assert schema.metadata
+
+ partition = schema.metadata.get(PARTITION_ON)
+ f = ddf.to_parquet(
+ compute=False,
+ path=tmp_path,
+ partition_on=partition,
+ engine="pyarrow",
+ overwrite=True,
+ write_metadata_file=True,
+ compression="brotli",
+ )
+ client.compute(f, sync=True, priority=2, resources=client_resources)
+ assert not os.path.exists(
+ self.path.as_posix()
+ ), f"already exits!: {self.path.as_posix()}"
+
+ if platform == "darwin":
+ subprocess.call(["mv", tmp_path.as_posix(), self.path.as_posix()])
+ else:
+ # -T will (should) cause the mv to fail if `path` wasn't successfully deleted
+ subprocess.call(["mv", "-T", tmp_path.as_posix(), self.path.as_posix()])
+ return True
+
+ def to_archive_symlink(
+ self,
+ client: Client,
+ ddf: dd.DataFrame,
+ is_partial: bool = False,
+ client_resources=None,
+ validate_after=True,
+ ) -> bool:
+ """
+ This differs from to_archive():
+ 1) to_parquet is run in this process. If the df is already
+ computed, there is no point in sending it to another worker
+ to write.
+
+ 2) symlink to next_numbered_path is created whether or not
+ is_partial (to_archive only does this on partials)
+
+ 3) we do not validate the written file. seems not useful to do
+ this, as the file will probably get overwritten on the next
+ loop anyway
+ """
+ path = self.partial_path if is_partial else self.path
+ next_numbered_path = self.next_numbered_path(path)
+ collection = self._collection
+ LOG.warning(f"{collection.merge_type.value}.to_archive_symlink()")
+
+ assert isinstance(ddf, dd.DataFrame), "must pass a dask df"
+
+ # We should validate before or after!!!
+ # _validate_df(self.compute(ddf), coll._schema)
+ target = (
+ next_numbered_path.name
+ ) # this is the symlink's target. it is a relative path (only the name)
+
+ schema = self._collection._schema
+ partition = schema.metadata.get(PARTITION_ON, None)
+ f = ddf.to_parquet(
+ compute=False,
+ path=next_numbered_path.as_posix(),
+ partition_on=partition,
+ engine="pyarrow",
+ overwrite=True,
+ write_metadata_file=True,
+ compression="brotli",
+ )
+ client.compute(f, sync=True, priority=2, resources=client_resources)
+
+ if os.path.exists(path.as_posix()) and not os.path.islink(path.as_posix()):
+ # This will fail when going from the old way to using symlinks,
+ # if self.path already exists and is a directory.
+ raise ValueError(
+ f"first time we run this, make sure the path doesnt exist: {path.as_posix()}"
+ )
+
+ if platform == "darwin":
+ subprocess.call(["ln", "-sfn", target, path.as_posix()])
+ else:
+ subprocess.call(["ln", "-sfnT", target, path.as_posix()])
+
+ if validate_after and not self.valid_archive(self.path):
+ LOG.error(f"{collection.merge_type.value} failed validation: {self.path}")
+ self.delete_archive(self.path)
+ return False
+ return True
+
+ # todo: unclear what the common interface should be here ... ?
+ def fetch(self, *args, **kwargs) -> pd.DataFrame | dd.DataFrame:
+ raise NotImplementedError("implement in subclass")
+
+ def build(self, *args, **kwargs) -> pd.DataFrame | dd.DataFrame:
+ raise NotImplementedError("implement in subclass")
+
+
+class MergeCollection(CollectionBase):
+ """Mergers take instances of DFCollections, and/or other Mergers"""
+
+ # In a merge, we can set offset = None which indicates that there is only 1
+ # period/item where the range is 'start' until now.
+ offset: str | None = Field(default="72h")
+ # In a merge, we can set start = None which indicates that there is only 1
+ # period/item where the range is (now - offset) until now.
+ start: AwareDatetimeISO | None = Field(
+ default=None,
+ description="This is the starting point in which data will"
+ " be retrieved in chunks from.",
+ frozen=True,
+ )
+
+ merge_type: MergeType | None = Field(default=None)
+ group_by: str | None = Field(default=None)
+ grouped_key: str | None = Field(default=None)
+ collection_item_class: type[MergeCollectionItem] = MergeCollectionItem
+
+ @model_validator(mode="after")
+ def check_start_and_offset_nullable(self) -> Self:
+ if self.offset is None and self.start is None:
+ raise AssertionError("cannot set both start and offset to None")
+ return self
+
+ @field_validator("merge_type")
+ def check_merge_type(cls, merge_type: MergeType | None, info: ValidationInfo):
+ if merge_type is None:
+ raise ValueError("Must explicitly provide a merge_type")
+
+ if merge_type not in MergeTypeSchemas:
+ raise ValueError("Must provide a supported merge_type")
+
+ return merge_type
+
+ # --- Properties ---
+ @property
+ def interval_start(self) -> datetime | None:
+ # if self.start is None and self.offset is set, the inferred start is (now - offset)
+ if self.start is None:
+ return datetime.now(tz=UTC).replace(microsecond=0) - pd.Timedelta(
+ self.offset
+ )
+ return self.start
+
+ @property
+ def items(self) -> list[MergeCollectionItem]:
+ items = []
+ for iv in self.interval_range:
+ cm = self.collection_item_class(start=iv[0])
+ cm._collection = self
+ items.append(cm)
+ return items
+
+ @property
+ def _schema(self) -> DataFrameSchema:
+ return MergeTypeSchemas[self.merge_type]
+
+ def signature(self) -> str:
+ arr = [
+ 1 if i.has_archive(include_empty=True) else 0
+ for i in self.items
+ if i.should_archive()
+ ]
+ repr_str = (
+ f"path={self.archive_path.as_posix()}; "
+ f"items={len(self.items)}; start={self.start} @ {self.offset}; {int(sum(arr) / len(arr) * 100)}% "
+ f"archived"
+ )
+ res = f"{self.__repr_name__()}({repr_str})"
+ return res
diff --git a/generalresearch/incite/mergers/foundations/enriched_session.py b/generalresearch/incite/mergers/foundations/enriched_session.py
index 049b1bc..4a300fc 100644
--- a/generalresearch/incite/mergers/foundations/enriched_session.py
+++ b/generalresearch/incite/mergers/foundations/enriched_session.py
@@ -14,7 +14,7 @@ from generalresearch.incite.collections.thl_web import (
SessionDFCollection,
WallDFCollection,
)
-from generalresearch.incite.mergers import (
+from generalresearch.incite.mergers.base import (
MergeCollection,
MergeCollectionItem,
MergeType,
diff --git a/generalresearch/incite/mergers/foundations/enriched_task_adjust.py b/generalresearch/incite/mergers/foundations/enriched_task_adjust.py
index e8a3654..f7c679f 100644
--- a/generalresearch/incite/mergers/foundations/enriched_task_adjust.py
+++ b/generalresearch/incite/mergers/foundations/enriched_task_adjust.py
@@ -12,7 +12,7 @@ from generalresearch.incite.collections.thl_web import (
TaskAdjustmentDFCollection,
)
from generalresearch.incite.exceptions import BuildError, BuildItemsError
-from generalresearch.incite.mergers import (
+from generalresearch.incite.mergers.base import (
MergeCollection,
MergeCollectionItem,
MergeType,
diff --git a/generalresearch/incite/mergers/foundations/enriched_wall.py b/generalresearch/incite/mergers/foundations/enriched_wall.py
index a74a556..9db5328 100644
--- a/generalresearch/incite/mergers/foundations/enriched_wall.py
+++ b/generalresearch/incite/mergers/foundations/enriched_wall.py
@@ -12,7 +12,7 @@ from generalresearch.incite.collections.thl_web import (
SessionDFCollection,
WallDFCollection,
)
-from generalresearch.incite.mergers import (
+from generalresearch.incite.mergers.base import (
MergeCollection,
MergeCollectionItem,
MergeType,
@@ -25,11 +25,11 @@ from generalresearch.incite.schemas.mergers.foundations.enriched_wall import (
EnrichedWallSchema,
)
from generalresearch.models.custom_types import UUIDStr
-from generalresearch.models.thl.user import User
from generalresearch.pg_helper import PostgresConfig
if TYPE_CHECKING:
from generalresearch.models.admin.request import ReportRequest
+ from generalresearch.models.thl.user import User
LOG = logging.getLogger("incite")
diff --git a/generalresearch/incite/mergers/foundations/user_id_product.py b/generalresearch/incite/mergers/foundations/user_id_product.py
index 863a741..3fb7b36 100644
--- a/generalresearch/incite/mergers/foundations/user_id_product.py
+++ b/generalresearch/incite/mergers/foundations/user_id_product.py
@@ -6,7 +6,7 @@ from typing import Any, Literal
from distributed import Client
from generalresearch.incite.collections.thl_web import UserDFCollection
-from generalresearch.incite.mergers import (
+from generalresearch.incite.mergers.base import (
MergeCollection,
MergeCollectionItem,
MergeType,
diff --git a/generalresearch/incite/mergers/pop_ledger.py b/generalresearch/incite/mergers/pop_ledger.py
index b32503c..cf63cad 100644
--- a/generalresearch/incite/mergers/pop_ledger.py
+++ b/generalresearch/incite/mergers/pop_ledger.py
@@ -9,7 +9,7 @@ from distributed import Client
from more_itertools import flatten
from generalresearch.incite.collections.thl_web import LedgerDFCollection
-from generalresearch.incite.mergers import (
+from generalresearch.incite.mergers.base import (
MergeCollection,
MergeCollectionItem,
MergeType,
diff --git a/generalresearch/incite/mergers/ym_survey_wall.py b/generalresearch/incite/mergers/ym_survey_wall.py
index a99e8ec..3cbf543 100644
--- a/generalresearch/incite/mergers/ym_survey_wall.py
+++ b/generalresearch/incite/mergers/ym_survey_wall.py
@@ -11,7 +11,7 @@ from sentry_sdk import capture_exception
from generalresearch.incite.collections.thl_web import WallDFCollection
from generalresearch.incite.exceptions import BuildError
-from generalresearch.incite.mergers import (
+from generalresearch.incite.mergers.base import (
MergeCollection,
MergeCollectionItem,
MergeType,
diff --git a/generalresearch/incite/mergers/ym_wall_summary.py b/generalresearch/incite/mergers/ym_wall_summary.py
index 69ef5c5..a01443b 100644
--- a/generalresearch/incite/mergers/ym_wall_summary.py
+++ b/generalresearch/incite/mergers/ym_wall_summary.py
@@ -13,7 +13,7 @@ from generalresearch.incite.collections.thl_web import (
WallDFCollection,
)
from generalresearch.incite.exceptions import FetchError
-from generalresearch.incite.mergers import (
+from generalresearch.incite.mergers.base import (
MergeCollection,
MergeCollectionItem,
MergeType,
diff --git a/generalresearch/incite/schemas/thl_web.py b/generalresearch/incite/schemas/thl_web.py
index c1be202..30c7076 100644
--- a/generalresearch/incite/schemas/thl_web.py
+++ b/generalresearch/incite/schemas/thl_web.py
@@ -1,6 +1,7 @@
from datetime import UTC, datetime, timedelta
import pandas as pd
+from grip_client.enums import AccessType
from pandera.pandas import Check, Column, DataFrameSchema, Index, MultiIndex
from generalresearch.incite.schemas import ARCHIVE_AFTER, ORDER_KEY
@@ -16,7 +17,6 @@ from generalresearch.models.thl.definitions import (
WallStatusCode2,
)
from generalresearch.models.thl.ledger import TransactionMetadataColumns
-from generalresearch.models.thl.maxmind.definitions import UserType
IP_REGEX_PATTERN = (
r"^((([0-9]|[1-9][0-9]|1[0-9]{2}|2[0-4][0-9]|25[0-5])\.){3}([0-9]|[1-9][0-9]|1[0-9]{2}|2[0-4]["
@@ -392,7 +392,7 @@ THLIPInfoSchema = DataFrameSchema(
dtype=str,
checks=[
Check.str_length(min_value=3, max_value=255),
- Check.isin([e.value for e in UserType]),
+ Check.isin([e.value for e in AccessType]),
],
nullable=True,
),
diff --git a/generalresearch/models/network/label.py b/generalresearch/models/network/label.py
index 60a6e58..b8fe4b0 100644
--- a/generalresearch/models/network/label.py
+++ b/generalresearch/models/network/label.py
@@ -51,6 +51,7 @@ class IPLabelSource(StrEnum):
INTERNAL_USE = "internal_use"
# An external "security" service flagged this IP
+ GRIP = "grip"
SPUR = "spur"
IPINFO = "ipinfo"
MAXMIND = "maxmind"
diff --git a/generalresearch/models/thl/definitions.py b/generalresearch/models/thl/definitions.py
index 0217a80..c40df21 100644
--- a/generalresearch/models/thl/definitions.py
+++ b/generalresearch/models/thl/definitions.py
@@ -185,7 +185,7 @@ class SessionStatusCode2(IntEnum, metaclass=ReprEnumMeta):
# Unable to parse either the bucket_id, request_id, or nudge_id from the url
ENTRY_URL_MODIFICATION = 1
- # The client's IP failed maxmind lookup, or we failed to store it for some reason
+ # The client's IP failed GRIP lookup, or we failed to store it for some reason
UNRECOGNIZED_IP = 2
# User is using an anonymous IP
USER_IS_ANONYMOUS = 3
diff --git a/generalresearch/models/thl/ipinfo.py b/generalresearch/models/thl/ipinfo.py
index e327bae..1e2be5b 100644
--- a/generalresearch/models/thl/ipinfo.py
+++ b/generalresearch/models/thl/ipinfo.py
@@ -5,6 +5,7 @@ from datetime import UTC, datetime
from typing import Any, Literal, Self
from faker import Faker
+from grip_client.enums import AccessType
from pydantic import (
BaseModel,
ConfigDict,
@@ -19,7 +20,6 @@ from generalresearch.models.custom_types import (
CountryISOLike,
IPvAnyAddressStr,
)
-from generalresearch.models.thl.maxmind.definitions import UserType
from generalresearch.pg_helper import PostgresConfig
fake = Faker()
@@ -173,11 +173,11 @@ class IPInformation(BaseModel):
default=None,
description="A score indicating the likelihood that the IP address is static.",
)
- user_type: UserType | None = Field(
+ user_type: AccessType | None = Field(
default=None,
description="The type of user associated with the IP address "
"(e.g., 'residential', 'business').",
- examples=[UserType.SCHOOL],
+ examples=[AccessType.RESIDENTIAL],
)
postal_code: str | None = Field(
default=None,
@@ -218,8 +218,8 @@ class IPInformation(BaseModel):
@property
def basic(self) -> bool:
- # This could be almost any field, but we're checking here if maxmind
- # insights was run on this record. If not, then most of the optional
+ # This could be almost any field, but we're checking here if GRIP
+ # was run on this record. If not, then most of the optional
# fields will be None
return self.is_anonymous is None
@@ -253,7 +253,7 @@ class IPInformation(BaseModel):
return d
@classmethod
- def from_mysql(cls, d: dict) -> Self:
+ def from_mysql(cls, d: dict[str, Any]) -> Self:
d["updated"] = d["updated"].replace(tzinfo=UTC)
return cls.model_validate(d)
@@ -261,3 +261,13 @@ class IPInformation(BaseModel):
class GeoIPInformation(IPInformation, IPGeoname):
model_config = ConfigDict(extra="ignore")
+
+ geoname_id: PositiveInt # type: ignore[reportIncompatibleVariableOverride]
+
+ @field_validator("geoname_id", mode="before")
+ @classmethod
+ def _coerce_geoname_id(cls, v: PositiveInt | None):
+ if v is None:
+ raise ValueError("GeoIPInformation can't be constructed")
+
+ return v
diff --git a/generalresearch/models/thl/ledger.py b/generalresearch/models/thl/ledger.py
index a9fbbb1..2b25d2e 100644
--- a/generalresearch/models/thl/ledger.py
+++ b/generalresearch/models/thl/ledger.py
@@ -22,12 +22,6 @@ from generalresearch.models.custom_types import (
UUIDStr,
check_valid_uuid,
)
-from generalresearch.models.thl.ledger_example import (
- _example_user_tx_adjustment,
- _example_user_tx_bonus,
- _example_user_tx_complete,
- _example_user_tx_payout,
-)
from generalresearch.models.thl.pagination import Page
from generalresearch.models.thl.payout_format import (
PayoutFormatType,
@@ -36,6 +30,53 @@ from generalresearch.models.thl.payout_format import (
from generalresearch.utils.enum import ReprEnumMeta
+def _example_user_tx_payout(schema: dict[str, Any]) -> None:
+
+ schema["example"] = UserLedgerTransactionUserPayout(
+ product_id=uuid4().hex,
+ payout_id=uuid4().hex,
+ amount=-5,
+ description="HIT Reward",
+ payout_format="${payout/100:.2f}",
+ created=datetime.now(tz=UTC),
+ ).model_dump(mode="json")
+
+
+def _example_user_tx_bonus(schema: dict[str, Any]) -> None:
+
+ schema["example"] = UserLedgerTransactionUserBonus(
+ product_id=uuid4().hex,
+ amount=100,
+ description="Compensation Bonus",
+ payout_format="${payout/100:.2f}",
+ created=datetime.now(tz=UTC),
+ ).model_dump(mode="json")
+
+
+def _example_user_tx_complete(schema: dict[str, Any]) -> None:
+
+ schema["example"] = UserLedgerTransactionTaskComplete(
+ product_id=uuid4().hex,
+ amount=38,
+ description="Task Complete",
+ payout_format="${payout/100:.2f}",
+ created=datetime.now(tz=UTC),
+ tsid=uuid4().hex,
+ ).model_dump(mode="json")
+
+
+def _example_user_tx_adjustment(schema: dict[str, Any]) -> None:
+
+ schema["example"] = UserLedgerTransactionTaskAdjustment(
+ product_id=uuid4().hex,
+ amount=-38,
+ description="Task Adjustment",
+ payout_format="${payout/100:.2f}",
+ created=datetime.now(tz=UTC),
+ tsid=uuid4().hex,
+ ).model_dump(mode="json")
+
+
class Direction(IntEnum, metaclass=ReprEnumMeta):
"""Entries on the debit side will increase debit normal accounts, while
entries on the credit side will decrease them. Conversely, entries on
@@ -393,7 +434,7 @@ class UserLedgerTransaction(BaseModel):
# It is optional b/c we'll calculate this from the query
balance_after: int | None = Field(default=None)
- def create_url(self, product_id: str):
+ def create_url(self, product_id: str) -> str | None:
raise NotImplementedError()
@computed_field(
@@ -431,7 +472,7 @@ class UserLedgerTransactionUserPayout(UserLedgerTransaction):
examples=["a3848e0a53d64f68a74ced5f61b6eb68"],
)
- def create_url(self, product_id: str):
+ def create_url(self, product_id: str) -> str | None:
return f"https://fsb.generalresearch.com/{product_id}/cashout/{self.payout_id}/"
@model_validator(mode="after")
@@ -459,7 +500,7 @@ class UserLedgerTransactionUserBonus(UserLedgerTransaction):
default="Compensation Bonus",
)
- def create_url(self, product_id: str):
+ def create_url(self, product_id: str) -> str | None:
return None
@model_validator(mode="after")
@@ -497,7 +538,7 @@ class UserLedgerTransactionTaskComplete(UserLedgerTransaction):
examples=["a3848e0a53d64f68a74ced5f61b6eb68"],
)
- def create_url(self, product_id: str):
+ def create_url(self, product_id: str) -> str | None:
return f"https://fsb.generalresearch.com/{product_id}/status/{self.tsid}/"
@model_validator(mode="after")
@@ -528,7 +569,7 @@ class UserLedgerTransactionTaskAdjustment(UserLedgerTransaction):
examples=["a3848e0a53d64f68a74ced5f61b6eb68"],
)
- def create_url(self, product_id: str):
+ def create_url(self, product_id: str) -> str | None:
return f"https://fsb.generalresearch.com/{product_id}/status/{self.tsid}/"
diff --git a/generalresearch/models/thl/ledger_example.py b/generalresearch/models/thl/ledger_example.py
deleted file mode 100644
index 0291691..0000000
--- a/generalresearch/models/thl/ledger_example.py
+++ /dev/null
@@ -1,64 +0,0 @@
-from __future__ import annotations
-
-from datetime import UTC, datetime
-from typing import Any
-from uuid import uuid4
-
-
-def _example_user_tx_payout(schema: dict[str, Any]) -> None:
- from generalresearch.models.thl.ledger import (
- UserLedgerTransactionUserPayout,
- )
-
- schema["example"] = UserLedgerTransactionUserPayout(
- product_id=uuid4().hex,
- payout_id=uuid4().hex,
- amount=-5,
- description="HIT Reward",
- payout_format="${payout/100:.2f}",
- created=datetime.now(tz=UTC),
- ).model_dump(mode="json")
-
-
-def _example_user_tx_bonus(schema: dict[str, Any]) -> None:
- from generalresearch.models.thl.ledger import (
- UserLedgerTransactionUserBonus,
- )
-
- schema["example"] = UserLedgerTransactionUserBonus(
- product_id=uuid4().hex,
- amount=100,
- description="Compensation Bonus",
- payout_format="${payout/100:.2f}",
- created=datetime.now(tz=UTC),
- ).model_dump(mode="json")
-
-
-def _example_user_tx_complete(schema: dict[str, Any]) -> None:
- from generalresearch.models.thl.ledger import (
- UserLedgerTransactionTaskComplete,
- )
-
- schema["example"] = UserLedgerTransactionTaskComplete(
- product_id=uuid4().hex,
- amount=38,
- description="Task Complete",
- payout_format="${payout/100:.2f}",
- created=datetime.now(tz=UTC),
- tsid=uuid4().hex,
- ).model_dump(mode="json")
-
-
-def _example_user_tx_adjustment(schema: dict[str, Any]) -> None:
- from generalresearch.models.thl.ledger import (
- UserLedgerTransactionTaskAdjustment,
- )
-
- schema["example"] = UserLedgerTransactionTaskAdjustment(
- product_id=uuid4().hex,
- amount=-38,
- description="Task Adjustment",
- payout_format="${payout/100:.2f}",
- created=datetime.now(tz=UTC),
- tsid=uuid4().hex,
- ).model_dump(mode="json")
diff --git a/generalresearch/models/thl/maxmind/__init__.py b/generalresearch/models/thl/maxmind/__init__.py
deleted file mode 100644
index e69de29..0000000
--- a/generalresearch/models/thl/maxmind/__init__.py
+++ /dev/null
diff --git a/generalresearch/models/thl/maxmind/definitions.py b/generalresearch/models/thl/maxmind/definitions.py
deleted file mode 100644
index 01431c7..0000000
--- a/generalresearch/models/thl/maxmind/definitions.py
+++ /dev/null
@@ -1,22 +0,0 @@
-from enum import Enum
-
-from generalresearch.utils.enum import ReprEnumMeta
-
-
-class UserType(Enum, metaclass=ReprEnumMeta):
- # https://support.maxmind.com/hc/en-us/articles/4408430082971-IP-Trait-Risk-Data#h_01FN6V8JMQMWZGWNPPAW77ZPY4
- BUSINESS = "business"
- CAFE = "cafe"
- CELLULAR = "cellular"
- COLLEGE = "college"
- CDN = "content_delivery_network"
- CPN = "consumer_privacy_network"
- GOVERNMENT = "government"
- HOSTING = "hosting"
- LIBRARY = "library"
- MILITARY = "military"
- RESIDENTIAL = "residential"
- ROUTER = "router"
- SCHOOL = "school"
- SEARCH_ENGINE = "search_engine_spider"
- TRAVELER = "traveler"
diff --git a/generalresearch/models/thl/session.py b/generalresearch/models/thl/session.py
index e4e264f..d547cf4 100644
--- a/generalresearch/models/thl/session.py
+++ b/generalresearch/models/thl/session.py
@@ -296,9 +296,9 @@ class WallBase(BaseModel):
finished: datetime | None = None,
) -> None:
# This should be called by the wall manager in order to actually update db
- from generalresearch import wall_status_codes
+ from generalresearch.wall_status_codes import annotate_status_code
- status, status_code_1, status_code_2 = wall_status_codes.annotate_status_code(
+ status, status_code_1, status_code_2 = annotate_status_code(
self.source,
ext_status_code_1,
ext_status_code_2,
@@ -317,18 +317,18 @@ class WallBase(BaseModel):
)
def is_soft_fail(self) -> bool:
- from generalresearch import wall_status_codes
+ from generalresearch.wall_status_codes import is_soft_fail
assert self.status is not None, "status should not be None"
assert self.status_code_1 is not None, "status_code_1 should not be None"
- return wall_status_codes.is_soft_fail(self)
+ return is_soft_fail(self)
def stop_marketplace_session(self) -> bool:
- from generalresearch import wall_status_codes
+ from generalresearch.wall_status_codes import stop_marketplace_session
assert self.status is not None, "status should not be None"
assert self.status_code_1 is not None, "status_code_1 should not be None"
- return wall_status_codes.stop_marketplace_session(self)
+ return stop_marketplace_session(self)
def get_status_after_adjustment(self) -> Status:
if self.adjusted_status in {
@@ -349,10 +349,13 @@ class WallBase(BaseModel):
WallAdjustedStatus.CPI_ADJUSTMENT,
}:
return self.adjusted_cpi
+
elif self.adjusted_status == WallAdjustedStatus.ADJUSTED_TO_FAIL:
return Decimal(0)
+
elif self.status == Status.COMPLETE:
return self.cpi
+
else:
return Decimal(0)
@@ -442,7 +445,7 @@ class Wall(WallBase):
d = self.model_dump(mode="json", exclude={"elapsed"})
return json.dumps(d)
- def model_dump_mysql(self, *args, **kwargs) -> dict:
+ def model_dump_mysql(self, *args, **kwargs) -> dict[str, Any]:
# Generate a dictionary representation of the model, with special handling for datetimes
d = self.model_dump(mode="json", exclude={"elapsed"}, *args, **kwargs)
d["started"] = self.started.replace(tzinfo=None)
@@ -496,13 +499,13 @@ class WallOut(WallBase):
)
# Serialize user_cpi to an int
- @field_serializer("user_cpi", return_type=int)
- def serialize_user_cpi(self, v: Decimal, _info):
+ @field_serializer("user_cpi", return_type=int | None)
+ def serialize_user_cpi(self, v: Decimal | None) -> int | None:
return decimal_to_int_cents(v)
# If user_cpi is an int, put it back to a decimal
@field_validator("user_cpi", mode="before")
- def deserialize_user_cpi(cls, v):
+ def deserialize_user_cpi(cls, v: Decimal | None) -> Decimal | None:
if isinstance(v, int):
return int_cents_to_decimal(v)
return v
@@ -510,7 +513,7 @@ class WallOut(WallBase):
# noinspection PyNestedDecorators
@field_validator("user_cpi", mode="after")
@classmethod
- def check_cpi_decimal_places(cls, v: Decimal) -> Decimal:
+ def check_cpi_decimal_places(cls, v: Decimal | None) -> Decimal | None:
if v is not None:
assert (
v.as_tuple().exponent >= -5
diff --git a/generalresearch/models/thl/user_iphistory.py b/generalresearch/models/thl/user_iphistory.py
index 257c41a..2773629 100644
--- a/generalresearch/models/thl/user_iphistory.py
+++ b/generalresearch/models/thl/user_iphistory.py
@@ -5,6 +5,7 @@ from datetime import UTC, datetime, timedelta
from typing import Self
from faker import Faker
+from grip_client.enums import AccessType
from pydantic import (
BaseModel,
ConfigDict,
@@ -22,7 +23,6 @@ from generalresearch.models.thl.ipinfo import (
GeoIPInformation,
normalize_ip,
)
-from generalresearch.models.thl.maxmind.definitions import UserType
from generalresearch.models.thl.user import User
from generalresearch.pg_helper import PostgresConfig
from generalresearch.redis_helper import RedisConfig
@@ -53,7 +53,11 @@ class UserIPRecord(BaseModel):
)
@property
- def user_type(self) -> UserType | None:
+ def user_type(self) -> AccessType | None:
+ return self.information.user_type if self.information else None
+
+ @property
+ def access_type(self) -> AccessType | None:
return self.information.user_type if self.information else None
@property
@@ -204,7 +208,6 @@ class UserIPHistory(BaseModel):
if res.get(x.ip):
x.information = res[x.ip]
-
def collapse_ip_records(self):
"""
- Records where sequential ipv6 addresses are in the same /64 block,
diff --git a/pyproject.toml b/pyproject.toml
index 03a1a1f..94f5073 100644
--- a/pyproject.toml
+++ b/pyproject.toml
@@ -59,4 +59,8 @@ addopts = "-v --tb=short"
target-version = "py314"
exclude = [
"generalresearch/thl_django",
-] \ No newline at end of file
+]
+
+[tool.pylint.messages_control]
+disable = ["all"]
+enable = ["cyclic-import"] \ No newline at end of file
diff --git a/test_utils/grliq/conftest.py b/test_utils/grliq/conftest.py
index 9a3bc56..891b73c 100644
--- a/test_utils/grliq/conftest.py
+++ b/test_utils/grliq/conftest.py
@@ -2,13 +2,13 @@ from __future__ import annotations
from collections.abc import Callable
from datetime import UTC, datetime, timedelta
+from typing import Any
from uuid import uuid4
import pytest
from pydantic import PostgresDsn
from generalresearch.config import GRLBaseSettings
-from generalresearch.grliq.managers import DUMMY_GRLIQ_DATA
from generalresearch.grliq.managers.forensic_data import (
GrlIqDataManager,
)
@@ -19,6 +19,10 @@ from generalresearch.grliq.managers.forensic_results import (
GrlIqCategoryResultsReader,
)
from generalresearch.grliq.models.forensic_data import GrlIqData
+from generalresearch.grliq.models.forensic_result import (
+ GrlIqCheckerResults,
+ GrlIqForensicCategoryResult,
+)
from generalresearch.pg_helper import PostgresConfig
# === Miscellaneous ===
@@ -75,11 +79,42 @@ def grliq_crr(grliq_db: PostgresConfig) -> GrlIqCategoryResultsReader:
# === Models ===
+@pytest.fixture(scope="session")
+def grliq_data_list() -> list[dict[str, Any]]:
+ return [
+ {
+ "data": GrlIqData.model_validate_json(
+ """{"mid": "3722ed29314940fabd37b42d808dcf5a", "uuid": "b11441da5a854dfbb8401d4c32e56db5", "phase": "offerwall-enter", "events": null, "vendor": "Google Inc.", "app_name": "Netscape", "calendar": "gregory", "language": "en-US", "platform": "Linux x86_64", "timezone": "America/Mexico_City", "client_ip": "131.196.250.250", "timestamp": "2025-02-27T16:05:34-06:00", "webrtc_ip": "131.196.250.250", "created_at": "2025-02-27T22:05:35.370589Z", "language_2": "en-US", "language_3": null, "platform_2": "Linux x86_64", "platform_3": null, "prefetched": true, "product_id": "d0606a0b5d034a8d81b1e3579d1f76fd", "webgl_flag": true, "webgl_hash": "da27e1b9b660057a3f5e185d3f5deabe", "canvas_hash": "14ed764326ec454d976c322261d99f16", "color_gamut": "3", "country_iso": "mx", "inner_width": 612, "outer_width": 1813, "product_sub": "20030107", "audio_codecs": "1,1,1,1,1,3,1,3,1,3,3,1,1,3,3,3,3,1,3,3,3,2,1,1", "cookie_check": "", "graphics_api": "WebKit WebGL", "inner_height": 1174, "mouse_events": null, "ontouchstart": false, "outer_height": 1261, "plugins_hash": "4c05fa2f766a444d4f253ead792c8b0e|2", "screen_width": 2560, "video_codecs": "1,3,3,3,3,3,3,3,3,3,1,1,1,1,1,1,3,1,1,1,3,3,1", "webgl_hash_2": "fc73fd5db75e2c36222fe34251be3971", "webrtc_error": false, "window_opera": false, "battery_level": 0.9, "canvas_hash_2": "bd11ebbf5c26fd20e0217820b4159752", "dynamic_range": false, "error_message": "Cannot read", "forced_colors": false, "math_result_1": "1.9275814160560204e-50", "math_result_2": "1.6182817135715877", "screen_height": 1440, "webgl_check_1": true, "webgl_context": "webgl2", "window_chrome": true, "connection_rtt": 150, "history_length": 16, "user_agent_str": "Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/131.0.0.0 Safari/537.36", "web_sql_exists": false, "calender_locale": "en-US", "connection_type": "", "inverted_colors": true, "navigator_brave": false, "product_user_id": "d1d55df1-959e-4740-b77c-fa1f4fc457ae", "request_headers": {"host": "test", "accept": "*/*", "connection": "keep-alive", "user-agent": "python-httpx/0.27.0", "content-length": "3646", "accept-encoding": "gzip, deflate", "x-forwarded-for": "131.196.250.250"}, "timezone_offset": 360, "webrtc_local_ip": "50486637-6b64-4812-b10a-0a75337c31bd.local", "battery_charging": true, "client_ip_detail": {"continent_code": "EU", "continent_name": "Europe", "country_name": "France", "is_in_european_union": true, "ip": "131.196.250.250", "isp": null, "latitude": null, "city_name": null, "longitude": null, "time_zone": null, "user_type": null, "country_iso": "mx", "postal_code": null, "is_anonymous": null, "accuracy_radius": null, "static_ip_score": null, "subdivision_1_iso": null, "subdivision_2_iso": null, "subdivision_1_name": null, "subdivision_2_name": null, "registered_country_iso": null}, "max_touch_points": 0, "numbering_system": "latn", "path_fingerprint": 3252, "prefers_contrast": "0", "rendering_engine": "WebKit", "timezone_success": "pass", "user_agent_hints": {"model": null, "brands": [{"brand": "Google Chrome", "version": "131"}, {"brand": "Chromium", "version": "131"}, {"brand": "Not_A Brand", "version": "24"}], "mobile": false, "bitness": "64", "platform": "Linux", "brands_full": [{"brand": "Google Chrome", "version": "131.0.6778.204"}, {"brand": "Chromium", "version": "131.0.6778.204"}, {"brand": "Not_A Brand", "version": "24.0.0.0"}], "architecture": "x86", "platform_version": "6.2.0"}, "user_agent_str_2": null, "webgl_extensions": "EXT_clip_control|EXT_color_buffer_float|EXT_color_buffer_half_float|EXT_conservative_depth|EXT_depth_clamp|EXT_disjoint_timer_query_webgl2|EXT_float_blend|EXT_polygon_offset_clamp|EXT_render_snorm|EXT_texture_compression_bptc|EXT_texture_compression_rgtc|EXT_texture_filter_anisotropic|EXT_texture_mirror_clamp_to_edge|EXT_texture_norm16|KHR_parallel_shader_compile|NV_shader_noperspective_interpolation|OES_draw_buffers_indexed|OES_sample_variables|OES_shader_multisample_interpolation|OES_texture_float_linear|OVR_multiview2|WEBGL_blend_func_extended|WEBGL_clip_cull_distance|WEBGL_compressed_texture_astc|WEBGL_compressed_texture_etc|WEBGL_compressed_texture_etc1|WEBGL_compressed_texture_s3tc|WEBGL_compressed_texture_s3tc_srgb|WEBGL_debug_renderer_info|WEBGL_debug_shaders|WEBGL_lose_context|WEBGL_multi_draw|WEBGL_polygon_mode|WEBGL_provoking_vertex|WEBGL_stencil_texturing", "webrtc_ip_detail": {"continent_code": "EU", "continent_name": "Europe", "country_name": "France", "is_in_european_union": true, "ip": "131.196.250.250", "isp": null, "latitude": null, "city_name": null, "longitude": null, "time_zone": null, "user_type": null, "country_iso": "mx", "postal_code": null, "is_anonymous": null, "accuracy_radius": null, "static_ip_score": null, "subdivision_1_iso": null, "subdivision_2_iso": null, "subdivision_1_name": null, "subdivision_2_name": null, "registered_country_iso": null}, "chrome_extensions": "", "execution_time_ms": 371.0999999642372, "graphics_renderer": "WebGL 2.0 (OpenGL ES 3.0 Chromium)", "keyboard_detected": true, "mime_types_length": 2, "request_fs_exists": true, "audio_context_flag": "pass", "audio_context_hash": "9307303774dec3248c18a939392090da", "canvas_fingerprint": 258, "canvas_pixel_check": false, "device_pixel_ratio": 1.0, "indexedDbData_blob": true, "navigator_keys_len": 79, "no_edge_pdf_plugin": false, "screen_avail_width": 2560, "webdriver_detected": false, "window_orientation": 0, "connection_downlink": 10.0, "navigator_webdriver": false, "non_native_function": false, "screen_avail_height": 1400, "supported_fonts_str": "72|768|262144|1073741824|0|0|540672|73728|7340032|1342177280|117446656|256|16|0|543|4290797636|1677723648|4168998400|0|1048576|262144|268500994|1342177280|262144|125829376|37888000|0|435363842|0|2147483648|109543424|1880099872|268435471", "text_2d_fingerprint": "bfcce91c9e71d11af7b14dbee4c75f83", "webrtc_is_supported": "pass", "canvas_support_level": "full", "do_not_track_enabled": "1", "hardware_concurrency": 12, "keyboard_layout_size": 48, "prefers_color_scheme": false, "webgl_max_anisotropy": 16, "battery_charging_time": 0.0, "browser_by_properties": "c", "eval_to_string_length": 33, "performance_loop_time": 0.09999996423721313, "session_storage_check": "pass", "unmasked_vendor_webgl": "Google Inc. (Intel)", "hardware_concurrency_2": 12, "hardware_concurrency_3": null, "localStorage_available": true, "memory_jsHeapSizeLimit": 4294705152, "mozilla_web_app_exists": false, "navigator_deviceMemory": 8.0, "navigator_java_enabled": false, "prefers_reduced_motion": false, "storage_estimate_quota": 1178717110272, "webdriver_detected_msg": "", "window_active_x_object": false, "window_external_exists": true, "color_depth_pixel_depth": "24-24", "indexedDbData_available": true, "navigator_cookieEnabled": true, "unmasked_renderer_webgl": "ANGLE (Intel, Mesa Intel(R) Graphics (RPL-P), OpenGL 4.6)", "battery_discharging_time": 0.0, "connection_effectiveType": "4g", "non_native_function_flag": "", "speech_synthesis_voice_1": "Google Bahasa Indonesia", "window_client_information": true, "audio_compressor_reduction": 20.538288116455078, "navigator_mediaDevices_len": 3, "audio_intensity_fingerprint": 124.04347527516074, "speech_synthesis_voice_hash": "8010ee3313813de521e48e63bd5a6f13", "microsoft_credentials_exists": false, "window_installTrigger_exists": false, "speech_synthesis_voices_count": 19, "webgl_shading_language_version": "WebGL GLSL ES 3.00 (OpenGL ES GLSL ES 3.0 Chromium)", "error_message_stack_access_count": 0, "speech_synthesis_avail_voices_count": 19, "error_message_stack_access_count_worker": 0}"""
+ ),
+ "result_data": GrlIqCheckerResults.model_validate_json(
+ """{"uuid": "b11441da5a854dfbb8401d4c32e56db5", "check_codecs": {"score": 0}, "check_timezone": {"score": 0}, "check_timestamp": {"score": 0}, "check_user_type": {"score": 0}, "check_ip_changes": {"score": 0}, "check_ip_country": {"score": 0}, "check_environment": {"score": 0}, "check_ip_timezone": {"score": 0}, "check_isp_changes": {"score": 0}, "check_useragent_js": {"score": 0}, "check_required_fonts": {"score": 0}, "check_user_anonymous": {"score": 0}, "check_webrtc_success": {"score": 0}, "check_seen_timestamps": {"msg": "duplicate timestamp", "score": 100}, "check_country_timezone": {"score": 0}, "check_prohibited_fonts": {"score": 0}, "check_timezone_changes": {"score": 0}, "check_execution_time_ms": {"msg": "duplicate execution_time_ms", "score": 100}, "check_fingerprint_reuse": {"score": 0}, "check_fingerprint_cycling": {"score": 0}, "check_ip_webrtc_ip_detail": {"score": 0}, "check_environment_critical": {"score": 0}, "check_useragent_other_enums": {"score": 0}, "check_useragent_ip_properties": {"score": 0}, "check_useragent_data_properties": {"score": 0}, "check_useragent_device_family_brand": {"score": 0}}"""
+ ),
+ "category_result": GrlIqForensicCategoryResult.model_validate_json(
+ """{"uuid": "b11441da5a854dfbb8401d4c32e56db5", "is_bot": 0, "is_tampered": 100, "is_velocity": 0, "is_anonymous": 0, "suspicious_ip": 0, "is_oscillating": 0, "is_teleporting": 0, "is_inconsistent": 0, "platform_ip_inconsistent": 0}"""
+ ),
+ "fraud_score": 100,
+ "is_attempt_allowed": False,
+ },
+ {
+ "data": GrlIqData.model_validate_json(
+ """{"mid": "35f6f5c30bc74ea7ac4aca7b40a02352", "uuid": "d54509f2f310499f8ab74839b10b2a41", "phase": "offerwall-enter", "events": null, "vendor": "Google Inc.", "app_name": "Netscape", "calendar": "gregory", "language": "en-US", "platform": "Linux x86_64", "timezone": "America/Los_Angeles", "client_ip": "104.9.125.144", "timestamp": "2025-02-28T11:34:39-08:00", "webrtc_ip": "172.56.209.195", "created_at": "2025-02-28T19:34:39.681872Z", "language_2": "en-US", "language_3": null, "platform_2": "Linux x86_64", "platform_3": null, "prefetched": true, "product_id": "d0606a0b5d034a8d81b1e3579d1f76fd", "webgl_flag": true, "webgl_hash": "da27e1b9b660057a3f5e185d3f5deabe", "canvas_hash": "e6e4d17da26050ce85ad00d3c6ea999e", "color_gamut": "3", "country_iso": "us", "inner_width": 841, "outer_width": 1680, "product_sub": "20030107", "audio_codecs": "1,1,1,1,1,3,1,3,1,3,3,1,1,3,3,3,3,1,3,3,3,2,1,1", "cookie_check": "", "graphics_api": "WebKit WebGL", "inner_height": 891, "mouse_events": null, "ontouchstart": false, "outer_height": 978, "plugins_hash": "4c05fa2f766a444d4f253ead792c8b0e|2", "screen_width": 1680, "video_codecs": "1,3,3,3,3,3,3,3,3,3,1,1,1,1,1,1,3,1,1,1,3,3,1", "webgl_hash_2": "fc73fd5db75e2c36222fe34251be3971", "webrtc_error": false, "window_opera": false, "battery_level": 0.41, "canvas_hash_2": "e0559d49b1864985cafc0d1c3a6b053c", "dynamic_range": false, "error_message": "Cannot read", "forced_colors": false, "math_result_1": "1.9275814160560204e-50", "math_result_2": "1.6182817135715877", "screen_height": 1050, "webgl_check_1": true, "webgl_context": "webgl2", "window_chrome": true, "connection_rtt": 100, "history_length": 11, "user_agent_str": "Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/131.0.0.0 Safari/537.36", "web_sql_exists": false, "calender_locale": "en-US", "connection_type": "", "inverted_colors": true, "navigator_brave": false, "product_user_id": "test-unit", "request_headers": {"dnt": "1", "host": "127.0.0.1:8081", "accept": "application/json, lk/null q=0.1", "origin": "http://127.0.0.1:8080", "referer": "http://127.0.0.1:8080/", "sec-ch-ua": "\\"Google Chrome\\";v=\\"131\\", \\"Chromium\\";v=\\"131\\", \\"Not_A Brand\\";v=\\"24\\"", "connection": "keep-alive", "user-agent": "Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/131.0.0.0 Safari/537.36", "content-type": "application/json", "content-length": "3313", "sec-fetch-dest": "empty", "sec-fetch-mode": "cors", "sec-fetch-site": "same-site", "accept-encoding": "gzip, deflate, br, zstd", "accept-language": "en-US,en;q=0.9", "sec-ch-ua-mobile": "?0", "sec-ch-ua-platform": "\\"Linux\\""}, "timezone_offset": 480, "webrtc_local_ip": "10.253.217.45,[2607:fb91:20c5:c6af:cda0:10b4:830a:a85e]", "battery_charging": false, "client_ip_detail": {"continent_code": "EU", "continent_name": "Europe", "country_name": "France", "is_in_european_union": true, "ip": "104.9.125.144", "isp": "AT&T Internet", "latitude": 37.3897, "city_name": "Mountain View", "longitude": -122.083, "time_zone": "America/Los_Angeles", "user_type": "residential", "country_iso": "us", "postal_code": "94041", "is_anonymous": false, "accuracy_radius": 5, "static_ip_score": 40.3, "subdivision_1_iso": "CA", "subdivision_2_iso": null, "subdivision_1_name": "California", "subdivision_2_name": null, "registered_country_iso": "us"}, "max_touch_points": 0, "numbering_system": "latn", "path_fingerprint": 3252, "prefers_contrast": "0", "rendering_engine": "WebKit", "timezone_success": "pass", "user_agent_hints": {"model": null, "brands": [{"brand": "Google Chrome", "version": "131"}, {"brand": "Chromium", "version": "131"}, {"brand": "Not_A Brand", "version": "24"}], "mobile": false, "bitness": "64", "platform": "Linux", "brands_full": [{"brand": "Google Chrome", "version": "131.0.6778.204"}, {"brand": "Chromium", "version": "131.0.6778.204"}, {"brand": "Not_A Brand", "version": "24.0.0.0"}], "architecture": "x86", "platform_version": "6.2.0"}, "user_agent_str_2": null, "webgl_extensions": "EXT_clip_control|EXT_color_buffer_float|EXT_color_buffer_half_float|EXT_conservative_depth|EXT_depth_clamp|EXT_disjoint_timer_query_webgl2|EXT_float_blend|EXT_polygon_offset_clamp|EXT_render_snorm|EXT_texture_compression_bptc|EXT_texture_compression_rgtc|EXT_texture_filter_anisotropic|EXT_texture_mirror_clamp_to_edge|EXT_texture_norm16|KHR_parallel_shader_compile|NV_shader_noperspective_interpolation|OES_draw_buffers_indexed|OES_sample_variables|OES_shader_multisample_interpolation|OES_texture_float_linear|OVR_multiview2|WEBGL_blend_func_extended|WEBGL_clip_cull_distance|WEBGL_compressed_texture_astc|WEBGL_compressed_texture_etc|WEBGL_compressed_texture_etc1|WEBGL_compressed_texture_s3tc|WEBGL_compressed_texture_s3tc_srgb|WEBGL_debug_renderer_info|WEBGL_debug_shaders|WEBGL_lose_context|WEBGL_multi_draw|WEBGL_polygon_mode|WEBGL_provoking_vertex|WEBGL_stencil_texturing", "webrtc_ip_detail": {"continent_code": "EU", "continent_name": "Europe", "country_name": "France", "is_in_european_union": true, "ip": "172.56.209.195", "isp": null, "latitude": null, "city_name": null, "longitude": null, "time_zone": null, "user_type": null, "country_iso": "us", "postal_code": null, "is_anonymous": null, "accuracy_radius": null, "static_ip_score": null, "subdivision_1_iso": null, "subdivision_2_iso": null, "subdivision_1_name": null, "subdivision_2_name": null, "registered_country_iso": null}, "chrome_extensions": "", "execution_time_ms": 924.5, "graphics_renderer": "WebGL 2.0 (OpenGL ES 3.0 Chromium)", "keyboard_detected": true, "mime_types_length": 2, "request_fs_exists": true, "audio_context_flag": "pass", "audio_context_hash": "9307303774dec3248c18a939392090da", "canvas_fingerprint": 258, "canvas_pixel_check": false, "device_pixel_ratio": 1.0, "indexedDbData_blob": true, "navigator_keys_len": 79, "no_edge_pdf_plugin": false, "screen_avail_width": 1680, "webdriver_detected": false, "window_orientation": 0, "connection_downlink": 10.0, "navigator_webdriver": false, "non_native_function": false, "screen_avail_height": 1010, "supported_fonts_str": "72|17152|327680|1073741824|0|0|540736|73728|7340032|1342177280|117446657|256|16|0|262687|4290797636|1677723648|4168998400|0|1048576|262144|268500994|1342177280|262144|125829376|37888000|0|435363842|0|2147483648|109543680|1880099888|301989903", "text_2d_fingerprint": "bfcce91c9e71d11af7b14dbee4c75f83", "webrtc_is_supported": "pass", "canvas_support_level": "full", "do_not_track_enabled": "1", "hardware_concurrency": 12, "keyboard_layout_size": 48, "prefers_color_scheme": false, "webgl_max_anisotropy": 16, "battery_charging_time": 0.0, "browser_by_properties": "c", "eval_to_string_length": 33, "performance_loop_time": 0.09999999962747097, "session_storage_check": "pass", "unmasked_vendor_webgl": "Google Inc. (Intel)", "hardware_concurrency_2": 12, "hardware_concurrency_3": null, "localStorage_available": true, "memory_jsHeapSizeLimit": 4294705152, "mozilla_web_app_exists": false, "navigator_deviceMemory": 8.0, "navigator_java_enabled": false, "prefers_reduced_motion": false, "storage_estimate_quota": 1178717110272, "webdriver_detected_msg": "", "window_active_x_object": false, "window_external_exists": true, "color_depth_pixel_depth": "24-24", "indexedDbData_available": true, "navigator_cookieEnabled": true, "unmasked_renderer_webgl": "ANGLE (Intel, Mesa Intel(R) Graphics (RPL-P), OpenGL 4.6)", "battery_discharging_time": 4844.0, "connection_effectiveType": "4g", "non_native_function_flag": "", "speech_synthesis_voice_1": "Google Bahasa Indonesia", "window_client_information": true, "audio_compressor_reduction": 20.538288116455078, "navigator_mediaDevices_len": 8, "audio_intensity_fingerprint": 124.04347527516074, "speech_synthesis_voice_hash": "8010ee3313813de521e48e63bd5a6f13", "microsoft_credentials_exists": false, "window_installTrigger_exists": false, "speech_synthesis_voices_count": 19, "webgl_shading_language_version": "WebGL GLSL ES 3.00 (OpenGL ES GLSL ES 3.0 Chromium)", "error_message_stack_access_count": 2, "speech_synthesis_avail_voices_count": 19, "error_message_stack_access_count_worker": 2}"""
+ ),
+ "result_data": GrlIqCheckerResults.model_validate_json(
+ """{"uuid": "d54509f2f310499f8ab74839b10b2a41", "check_codecs": {"score": 0}, "check_timezone": {"score": 0}, "check_timestamp": {"score": 0}, "check_user_type": {"score": 0}, "check_ip_changes": {"score": 0}, "check_ip_country": {"score": 0}, "check_environment": {"msg": "error_message_stack_access_count: 2", "score": 100}, "check_ip_timezone": {"score": 0}, "check_isp_changes": {"score": 0}, "check_useragent_js": {"score": 0}, "check_required_fonts": {"score": 0}, "check_user_anonymous": {"score": 0}, "check_webrtc_success": {"score": 0}, "check_seen_timestamps": {"score": 0}, "check_country_timezone": {"score": 0}, "check_prohibited_fonts": {"score": 0}, "check_timezone_changes": {"score": 0}, "check_execution_time_ms": {"score": 0}, "check_fingerprint_reuse": {"score": 0}, "check_fingerprint_cycling": {"score": 0}, "check_ip_webrtc_ip_detail": {"score": 0}, "check_environment_critical": {"score": 0}, "check_useragent_other_enums": {"score": 0}, "check_useragent_ip_properties": {"score": 0}, "check_useragent_data_properties": {"score": 0}, "check_useragent_device_family_brand": {"score": 0}}"""
+ ),
+ "category_result": GrlIqForensicCategoryResult.model_validate_json(
+ """{"uuid": "d54509f2f310499f8ab74839b10b2a41", "is_bot": 0, "is_tampered": 0, "is_velocity": 0, "is_anonymous": 0, "suspicious_ip": 0, "is_oscillating": 0, "is_teleporting": 0, "is_inconsistent": 10, "platform_ip_inconsistent": 0}"""
+ ),
+ "fraud_score": 10,
+ "is_attempt_allowed": True,
+ },
+ ]
+
+
@pytest.fixture(scope="function")
-def grliq_data() -> GrlIqData:
- from generalresearch.grliq.managers import DUMMY_GRLIQ_DATA
+def grliq_data(grliq_data_list: list[dict[str, Any]]) -> GrlIqData:
- g: GrlIqData = DUMMY_GRLIQ_DATA[1]["data"]
+ g: GrlIqData = grliq_data_list[1]["data"]
g.id = None
g.uuid = uuid4().hex
@@ -89,7 +124,9 @@ def grliq_data() -> GrlIqData:
@pytest.fixture
-def grliq_data_factory(grliq_dm: GrlIqDataManager) -> Callable[..., GrlIqData]:
+def grliq_data_factory(
+ grliq_dm: GrlIqDataManager, grliq_data_list: list[dict[str, Any]]
+) -> Callable[..., GrlIqData]:
def _inner(
is_attempt_allowed: bool = True,
@@ -109,9 +146,8 @@ def grliq_data_factory(grliq_dm: GrlIqDataManager) -> Callable[..., GrlIqData]:
:param mid: the thl_session:uuid / mid for the attempt.
:return:
"""
- import copy
- res: GrlIqData = copy.deepcopy(DUMMY_GRLIQ_DATA[int(is_attempt_allowed)])
+ res: GrlIqData = grliq_data_list[int(is_attempt_allowed)]["data"]
product_id = product_id or uuid4().hex
product_user_id = product_user_id or uuid4().hex
diff --git a/test_utils/incite/mergers/conftest.py b/test_utils/incite/mergers/conftest.py
index c0f0bcf..1f88804 100644
--- a/test_utils/incite/mergers/conftest.py
+++ b/test_utils/incite/mergers/conftest.py
@@ -6,7 +6,7 @@ from datetime import datetime, timedelta
import pytest
from generalresearch.incite.base import GRLDatasets
-from generalresearch.incite.mergers import MergeType
+from generalresearch.incite.mergers.base import MergeType
from generalresearch.incite.mergers.foundations.enriched_session import (
EnrichedSessionMerge,
)