diff options
| author | Max Nanis | 2026-08-29 22:00:43 -0700 |
|---|---|---|
| committer | Max Nanis | 2026-08-29 22:00:43 -0700 |
| commit | dda0067fdfe563e8270562f2e0c9da7c74c37a0b (patch) | |
| tree | 53e742f9c7b585b1a8e46e4fe4f1fd04c629d196 | |
| parent | 8219bf4814a7aa374a500516f2254f39085f0357 (diff) | |
| download | generalresearch-dda0067fdfe563e8270562f2e0c9da7c74c37a0b.tar.gz generalresearch-dda0067fdfe563e8270562f2e0c9da7c74c37a0b.zip | |
pylint ⭕️ cyclic-import checks
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, ) |
