| import os |
| import io |
| import uvicorn |
| from fastapi import FastAPI, Query, UploadFile, File |
| from typing import List, Optional, Dict, Any |
| import openpyxl |
| import xml.etree.ElementTree as ET |
| from fuzzywuzzy import fuzz |
| import pandas as pd |
| import numpy as np |
| from pydantic import BaseModel |
| import requests |
| from datetime import datetime |
| import psycopg2 |
| from psycopg2.extras import RealDictCursor |
|
|
| app = FastAPI(title="Search API", description="API for searching individuals and entities data") |
|
|
| |
| def parse_xml(): |
| tree = ET.parse('consolidated.xml') |
| return tree.getroot() |
|
|
| root = parse_xml() |
|
|
| |
| def extract_text(element, path): |
| found = element.find(path) |
| return found.text.strip() if found is not None and found.text else "" |
|
|
| def calculate_similarity_score(row, search_term): |
| """Given a row (iterable of values) and a single search term (string), |
| compute the max fuzz.partial_ratio.""" |
| scores = [] |
| for value in row: |
| if pd.isna(value) or not isinstance(value, str): |
| continue |
| ratio = fuzz.partial_ratio(search_term.lower(), value.lower()) |
| scores.append(ratio) |
| return max(scores) if scores else 0 |
|
|
| def get_individuals_data(): |
| """Load individual data from XML.""" |
| xml_root = parse_xml() |
| individuals = [] |
| for person in xml_root.findall('.//INDIVIDUAL'): |
| individual = { |
| 'DATAID': extract_text(person, './/DATAID'), |
| 'FIRST_NAME': extract_text(person, './/FIRST_NAME'), |
| 'SECOND_NAME': extract_text(person, './/SECOND_NAME'), |
| 'THIRD_NAME': extract_text(person, './/THIRD_NAME'), |
| 'NATIONALITY': extract_text(person, './/NATIONALITY/VALUE'), |
| 'GENDER': extract_text(person, './/GENDER'), |
| 'DOB': extract_text(person, './/INDIVIDUAL_DATE_OF_BIRTH/DATE'), |
| 'POB': extract_text(person, './/INDIVIDUAL_PLACE_OF_BIRTH/CITY'), |
| 'COUNTRY_OF_BIRTH': extract_text(person, './/INDIVIDUAL_PLACE_OF_BIRTH/COUNTRY'), |
| } |
| individuals.append(individual) |
| return pd.DataFrame(individuals) |
|
|
| def get_entities_data(): |
| """Load entity data from XML.""" |
| xml_root = parse_xml() |
| entities = [] |
| for entity in xml_root.findall('.//ENTITY'): |
| entity_data = { |
| 'DATAID': extract_text(entity, './/DATAID'), |
| 'FIRST_NAME': extract_text(entity, './/FIRST_NAME'), |
| 'ENTITY_ALIAS': extract_text(entity, './/ENTITY_ALIAS/ALIAS_NAME'), |
| 'ENTITY_ADDRESS': extract_text(entity, './/ENTITY_ADDRESS/CITY'), |
| 'ENTITY_COUNTRY': extract_text(entity, './/ENTITY_ADDRESS/COUNTRY'), |
| } |
| entities.append(entity_data) |
| return pd.DataFrame(entities) |
|
|
| class SearchResponse(BaseModel): |
| total_results: int |
| |
| confidence_threshold: Optional[float] = None |
| results: List[dict] |
|
|
| @app.get("/search/individuals", response_model=SearchResponse) |
| async def search_individuals( |
| query: str = Query(..., description="Search term"), |
| min_confidence: float = Query(50.0, description="Minimum confidence score (0-100)", ge=0, le=100) |
| ): |
| individuals_df = get_individuals_data() |
| |
| |
| similarity_scores = individuals_df.apply( |
| lambda x: calculate_similarity_score(x, query), |
| axis=1 |
| ) |
| |
| |
| mask = similarity_scores >= min_confidence |
| filtered_df = individuals_df[mask].copy() |
| |
| |
| filtered_df['match_confidence'] = similarity_scores[mask].round(1) |
| |
| |
| filtered_df = filtered_df.sort_values('match_confidence', ascending=False) |
| |
| |
| results = filtered_df.to_dict('records') |
| |
| print(f"Search complete. Total results: {len(results)}") |
| |
| return { |
| "total_results": len(results), |
| "results": results |
| } |
|
|
| @app.get("/search/entities", response_model=SearchResponse) |
| async def search_entities( |
| query: str = Query(..., description="Search term"), |
| min_confidence: float = Query(50.0, description="Minimum confidence score (0-100)", ge=0, le=100) |
| ): |
| entities_df = get_entities_data() |
| |
| |
| similarity_scores = entities_df.apply( |
| lambda x: calculate_similarity_score(x, query), |
| axis=1 |
| ) |
| |
| |
| mask = similarity_scores >= min_confidence |
| filtered_df = entities_df[mask].copy() |
| |
| |
| filtered_df['match_confidence'] = similarity_scores[mask].round(1) |
| |
| |
| filtered_df = filtered_df.sort_values('match_confidence', ascending=False) |
| |
| |
| results = filtered_df.to_dict('records') |
| |
| return { |
| "total_results": len(results), |
| "confidence_threshold": min_confidence, |
| "results": results |
| } |
|
|
| @app.post("/query") |
| async def query_endpoint( |
| llm_provider: str = "anthropic", |
| k: int = Query(3, description="Number of results to return"), |
| temperature: float = Query(0, description="Temperature for LLM"), |
| max_tokens: int = Query(100, description="Maximum tokens to generate"), |
| query: str = Query(..., description="Query string"), |
| llm_model: str = Query("claude-3-opus-latest", description="LLM model to use") |
| ): |
| url = "https://api.edenai.run/v2/aiproducts/askyoda/v2/3b3f2311-8c4b-41f7-8d13-ac78982e8311/query" |
| |
| payload = { |
| "llm_provider": llm_provider, |
| "k": k, |
| "temperature": temperature, |
| "max_tokens": max_tokens, |
| "query": query, |
| "llm_model": llm_model |
| } |
| |
| headers = { |
| "accept": "application/json", |
| "content-type": "application/json", |
| "authorization": "Bearer eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJ1c2VyX2lkIjoiM2NlMTljZDItNTk2Zi00NzRhLWE3OGMtZWZhNTM1YWJlNmY0IiwidHlwZSI6ImFwaV90b2tlbiJ9.x_PPBluoOy2GT3ZjDQ3dcS8WtbNn95RfkoFKkhM_a5A" |
| } |
| |
| response = requests.post(url, json=payload, headers=headers) |
| return response.json() |
|
|
| @app.post("/search_sanctions") |
| async def search_sanctions( |
| file: UploadFile = File(...), |
| min_confidence: float = Query(50.0, description="Minimum confidence score (0-100)", ge=0, le=100) |
| ): |
| print(f"Starting sanctions search process...") |
| |
| |
| print("Loading sanctioned lists...") |
| d1_df = pd.read_csv('D1.csv') |
| d2_df = pd.read_csv('D2.csv') |
| print(f"Loaded {len(d1_df)} individuals and {len(d2_df)} entities from sanctions lists") |
| |
| |
| print(f"Saving uploaded file as temp2.xlsx") |
| file_content = await file.read() |
| with open('temp2.xlsx', 'wb') as f: |
| f.write(file_content) |
| |
| |
| print(f"Reading search terms from temp2.xlsx") |
| search_df = pd.read_excel('temp2.xlsx') |
| print(f"Loaded {len(search_df)} search terms") |
| |
| all_matches = [] |
| |
| try: |
| conn = psycopg2.connect( |
| host="13.126.242.31", |
| database="aml", |
| user="dev_cbs_admin", |
| password="Finovate@2023" |
| ) |
| cur = conn.cursor() |
| |
| |
| for _, search_row in search_df.iterrows(): |
| |
| customer_name = search_row['Customer Name'] if 'Customer Name' in search_row else str(search_row.iloc[0]) |
| print(f"Processing customer: {customer_name[:100]}...") |
| |
| matches = [] |
| |
| |
| for _, individual in d1_df.iterrows(): |
| |
| sanction_name = ' '.join([ |
| str(individual['FIRST_NAME']) if pd.notna(individual['FIRST_NAME']) else '', |
| str(individual['SECOND_NAME']) if pd.notna(individual['SECOND_NAME']) else '', |
| str(individual['THIRD_NAME']) if pd.notna(individual['THIRD_NAME']) else '' |
| ]).strip() |
| |
| |
| score = calculate_similarity_score([customer_name], sanction_name) |
| |
| if score >= min_confidence: |
| |
| insert_query = """ |
| INSERT INTO sanctions_matches ( |
| search_term_customer_name, matched_list, matched_name, |
| confidence_score, dataid, first_name, second_name, |
| third_name, nationality, gender, dob, pob, country_of_birth, |
| status |
| ) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s) |
| """ |
| |
| values = ( |
| customer_name, |
| 'D1 - Individuals', |
| sanction_name, |
| round(score, 2), |
| individual.get('DATAID'), |
| individual.get('FIRST_NAME'), |
| individual.get('SECOND_NAME'), |
| individual.get('THIRD_NAME'), |
| individual.get('NATIONALITY'), |
| individual.get('GENDER'), |
| individual.get('DOB'), |
| individual.get('POB'), |
| individual.get('COUNTRY_OF_BIRTH'), |
| 'pending' |
| ) |
| |
| cur.execute(insert_query, values) |
| |
| |
| matched_data_dict = {k: (None if pd.isna(v) else v) for k, v in individual.to_dict().items()} |
| match = { |
| 'search_term_customer_name': customer_name, |
| 'matched_list': 'D1 - Individuals', |
| 'matched_name': sanction_name, |
| 'matched_data': matched_data_dict, |
| 'confidence_score': round(score, 2) |
| } |
| matches.append(match) |
| |
| |
| for _, entity in d2_df.iterrows(): |
| entity_name = ' '.join([ |
| str(entity['FIRST_NAME']) if pd.notna(entity['FIRST_NAME']) else '', |
| str(entity['ENTITY_ALIAS']) if pd.notna(entity['ENTITY_ALIAS']) else '' |
| ]).strip() |
| |
| |
| score = calculate_similarity_score([customer_name], entity_name) |
| |
| if score >= min_confidence: |
| |
| insert_query = """ |
| INSERT INTO sanctions_matches ( |
| search_term_customer_name, matched_list, matched_name, |
| confidence_score, dataid, first_name, entity_alias, |
| entity_address, entity_country, status |
| ) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s) |
| """ |
| |
| values = ( |
| customer_name, |
| 'D2 - Entities', |
| entity_name, |
| round(score, 2), |
| entity.get('DATAID'), |
| entity.get('FIRST_NAME'), |
| entity.get('ENTITY_ALIAS'), |
| entity.get('ENTITY_ADDRESS'), |
| entity.get('ENTITY_COUNTRY'), |
| 'pending' |
| ) |
| |
| cur.execute(insert_query, values) |
| |
| |
| matched_data_dict = {k: (None if pd.isna(v) else v) for k, v in entity.to_dict().items()} |
| match = { |
| 'search_term_customer_name': customer_name, |
| 'matched_list': 'D2 - Entities', |
| 'matched_name': entity_name, |
| 'matched_data': matched_data_dict, |
| 'confidence_score': round(score, 2) |
| } |
| matches.append(match) |
| |
| if matches: |
| |
| matches.sort(key=lambda x: x['confidence_score'], reverse=True) |
| all_matches.extend(matches) |
| |
| |
| conn.commit() |
| print(f"Response code: 200") |
| |
| return { |
| "status": "success", |
| "code": 200, |
| "message": "Successfully updated sanctions_matches table" |
| } |
| |
| except Exception as e: |
| print(f"Database error: {str(e)}") |
| if 'conn' in locals(): |
| conn.rollback() |
| print(f"Response code: 400") |
| return { |
| "status": "error", |
| "code": 400 |
| } |
| finally: |
| if 'cur' in locals(): |
| cur.close() |
| if 'conn' in locals(): |
| conn.close() |
|
|
| @app.post("/large_transactions") |
| async def large_transactions( |
| file: UploadFile = File(...), |
| limit: float = Query(1000000.0, description="Minimum transaction amount (default: 10,00,000)") |
| ): |
| print(f"Starting large transactions analysis with limit: {limit}...") |
| |
| |
| print(f"Reading transaction data from: {file.filename}") |
| file_content = await file.read() |
| excel_io = io.BytesIO(file_content) |
| df = pd.read_excel(excel_io) |
| |
| |
| large_txns = df[df['AmountINR'] > limit].copy() |
| print(f"Found {len(large_txns)} transactions above {limit}") |
| |
| |
| results = [] |
| for _, row in large_txns.iterrows(): |
| row_dict = {k: (None if pd.isna(v) else v) for k, v in row.to_dict().items()} |
| results.append(row_dict) |
|
|
| |
| try: |
| conn = psycopg2.connect( |
| host="13.126.242.31", |
| database="aml", |
| user="dev_cbs_admin", |
| password="Finovate@2023" |
| ) |
| cur = conn.cursor(cursor_factory=RealDictCursor) |
| |
| |
| for transaction in results: |
| |
| for key, value in transaction.items(): |
| if isinstance(value, datetime): |
| transaction[key] = value.isoformat() |
|
|
| |
| cur.execute(""" |
| SELECT column_name |
| FROM information_schema.columns |
| WHERE table_name = 'large_transactions' |
| """) |
| db_columns = [row['column_name'] for row in cur.fetchall()] |
| print("Database columns:", db_columns) |
|
|
| |
| cleaned_transaction = {} |
| for key, value in transaction.items(): |
| cleaned_key = key.replace('.', '_').replace(' ', '') |
| if cleaned_key.lower() in [col.lower() for col in db_columns]: |
| correct_case = next(col for col in db_columns if col.lower() == cleaned_key.lower()) |
| cleaned_transaction[correct_case] = value |
|
|
| print("Number of columns:", len(cleaned_transaction)) |
| print("Cleaned transaction keys:", cleaned_transaction.keys()) |
|
|
| |
| cleaned_transaction['status'] = 'pending' |
| |
| columns = ', '.join(cleaned_transaction.keys()) |
| values = ', '.join(['%s'] * len(cleaned_transaction)) |
| |
| insert_query = f""" |
| INSERT INTO large_transactions ({columns}) |
| VALUES ({values}) |
| """ |
| print("Insert query:", insert_query) |
|
|
| try: |
| values_to_insert = tuple(cleaned_transaction.values()) |
| cur.execute(insert_query, values_to_insert) |
| except psycopg2.Error as e: |
| return { |
| "status": "error", |
| "code": 400 |
| } |
| |
| try: |
| conn.commit() |
| return { |
| "status": "success", |
| "code": 200 |
| } |
| except Exception as e: |
| print(f"Commit error: {str(e)}") |
| return { |
| "status": "error", |
| "message": f"Failed to commit transactions: {str(e)}" |
| } |
| except Exception as e: |
| print(f"Database error: {str(e)}") |
| return { |
| "status": "error", |
| "message": f"Database error: {str(e)}" |
| } |
| finally: |
| if 'cur' in locals(): |
| cur.close() |
| if 'conn' in locals(): |
| conn.close() |
|
|
| @app.post("/reset_db") |
| async def reset_db(): |
| try: |
| conn = psycopg2.connect( |
| host="13.126.242.31", |
| database="aml", |
| user="dev_cbs_admin", |
| password="Finovate@2023" |
| ) |
| cur = conn.cursor(cursor_factory=RealDictCursor) |
| |
| |
| delete_large_txn_query = "DELETE FROM large_transactions" |
| delete_sanctions_query = "DELETE FROM sanctions_matches" |
| delete_velocity_query = "DELETE FROM velocity" |
| delete_backup_query = "DELETE FROM backup_data" |
| |
| cur.execute(delete_large_txn_query) |
| large_txn_rows = cur.rowcount |
| |
| cur.execute(delete_sanctions_query) |
| sanctions_rows = cur.rowcount |
| |
| cur.execute(delete_velocity_query) |
| velocity_rows = cur.rowcount |
| |
| cur.execute(delete_backup_query) |
| backup_rows = cur.rowcount |
| |
| try: |
| conn.commit() |
| return { |
| "status": "success", |
| "message": f"Successfully deleted {large_txn_rows} records from large_transactions table, {sanctions_rows} records from sanctions_matches table, {velocity_rows} records from velocity table, and {backup_rows} records from backup_data table" |
| } |
| except Exception as e: |
| print(f"Commit error: {str(e)}") |
| return { |
| "status": "error", |
| "message": f"Failed to commit deletion: {str(e)}" |
| } |
| |
| except Exception as e: |
| print(f"Database error: {str(e)}") |
| return { |
| "status": "error", |
| "message": f"Database error: {str(e)}" |
| } |
| finally: |
| if 'cur' in locals(): |
| cur.close() |
| if 'conn' in locals(): |
| conn.close() |
|
|
| @app.post("/sanction_list_2") |
| async def sanction_list_2( |
| file: UploadFile = File(...), |
| min_confidence: float = Query(50.0, description="Minimum confidence score (0-100)", ge=0, le=100) |
| ): |
| print("Starting sanctions list 2 matching process...") |
| |
| |
| print("Loading sanctioned lists...") |
| d1_df = pd.read_csv('D1.csv') |
| d2_df = pd.read_csv('D2.csv') |
| print(f"Loaded {len(d1_df)} individuals and {len(d2_df)} entities from sanctions lists") |
| |
| |
| print(f"Saving uploaded file as temp1.xlsx") |
| file_content = await file.read() |
| with open('temp1.xlsx', 'wb') as f: |
| f.write(file_content) |
| |
| |
| print(f"Reading transaction data from temp1.xlsx") |
| txn_df = pd.read_excel('temp1.xlsx') |
| print(f"Loaded {len(txn_df)} transactions") |
| |
| all_matches = [] |
| |
| def check_name_against_sanctions(name, source_field): |
| """Helper function to compare `name` against both D1 & D2.""" |
| matches = [] |
| |
| |
| for _, individual in d1_df.iterrows(): |
| sanction_name = ' '.join([ |
| str(individual['FIRST_NAME']) if pd.notna(individual['FIRST_NAME']) else '', |
| str(individual['SECOND_NAME']) if pd.notna(individual['SECOND_NAME']) else '', |
| str(individual['THIRD_NAME']) if pd.notna(individual['THIRD_NAME']) else '' |
| ]).strip() |
| |
| score = calculate_similarity_score([name], sanction_name) |
| if score >= min_confidence: |
| matched_data_dict = {k: (None if pd.isna(v) else v) for k, v in individual.to_dict().items()} |
| matches.append({ |
| 'source_field': source_field, |
| 'matched_name': name, |
| 'matched_list': 'D1 - Individuals', |
| 'matched_data': matched_data_dict, |
| 'confidence_score': round(score, 2) |
| }) |
| |
| |
| for _, entity in d2_df.iterrows(): |
| entity_name = ' '.join([ |
| str(entity['FIRST_NAME']) if pd.notna(entity['FIRST_NAME']) else '', |
| str(entity['ENTITY_ALIAS']) if pd.notna(entity['ENTITY_ALIAS']) else '' |
| ]).strip() |
| score = calculate_similarity_score([name], entity_name) |
| if score >= min_confidence: |
| matched_data_dict = {k: (None if pd.isna(v) else v) for k, v in entity.to_dict().items()} |
| matches.append({ |
| 'source_field': source_field, |
| 'matched_name': name, |
| 'matched_list': 'D2 - Entities', |
| 'matched_data': matched_data_dict, |
| 'confidence_score': round(score, 2) |
| }) |
| |
| return matches |
| |
| |
| for _, txn in txn_df.iterrows(): |
| |
| if pd.notna(txn.get('SenderName')): |
| matches = check_name_against_sanctions(txn['SenderName'], 'SenderName') |
| if matches: |
| txn_dict = {k: (None if pd.isna(v) else v) for k, v in txn.to_dict().items()} |
| for match in matches: |
| match['transaction_data'] = txn_dict |
| all_matches.extend(matches) |
| |
| |
| if pd.notna(txn.get('BeneficiaryName')): |
| matches = check_name_against_sanctions(txn['BeneficiaryName'], 'BeneficiaryName') |
| if matches: |
| txn_dict = {k: (None if pd.isna(v) else v) for k, v in txn.to_dict().items()} |
| for match in matches: |
| match['transaction_data'] = txn_dict |
| all_matches.extend(matches) |
| |
| |
| all_matches.sort(key=lambda x: x['confidence_score'], reverse=True) |
| |
| print(f"Search complete. Found {len(all_matches)} matches above confidence threshold") |
| |
| return { |
| "total_transactions": len(txn_df), |
| "total_matches": len(all_matches), |
| "min_confidence_threshold": min_confidence, |
| "matches": all_matches |
| } |
|
|
| @app.post("/velocity_analysis") |
| async def velocity_analysis( |
| file: UploadFile = File(...), |
| period: str = Query("M", description="Period for velocity calculation: D(days), W(weeks), M(months)"), |
| top_percent: float = Query(4.0, description="Top percentage to flag (default: 4%)", ge=0, le=100) |
| ): |
| print(f"Starting velocity analysis for period: {period}") |
| |
| |
| print(f"Saving uploaded file as temp1.xlsx") |
| file_content = await file.read() |
| with open('temp1.xlsx', 'wb') as f: |
| f.write(file_content) |
| |
| |
| print(f"Reading transaction data from temp1.xlsx") |
| df = pd.read_excel('temp1.xlsx') |
| |
| |
| if 'TransactionDate' not in df.columns: |
| return {"error": "TransactionDate column is required in the Excel file"} |
| |
| |
| df['TransactionDate'] = pd.to_datetime(df['TransactionDate']) |
| |
| |
| df = df.sort_values('TransactionDate') |
| |
| account_metrics = [] |
| for account in df['BeneficiaryAccountNumber'].unique(): |
| account_df = df[df['BeneficiaryAccountNumber'] == account] |
| |
| date_range = account_df['TransactionDate'].max() - account_df['TransactionDate'].min() |
| |
| |
| if period == 'D': |
| period_length = date_range.days |
| period_name = 'days' |
| elif period == 'W': |
| period_length = date_range.days / 7 |
| period_name = 'weeks' |
| else: |
| period_length = date_range.days / 30 |
| period_name = 'months' |
| |
| |
| num_transactions = len(account_df) |
| if period_length == 0: |
| period_length = 1 |
| velocity = num_transactions / period_length |
| |
| latest_record = account_df.iloc[-1] |
| |
| account_metrics.append({ |
| 'account_number': account, |
| 'beneficiary_name': latest_record.get('BeneficiaryName', 'N/A'), |
| 'total_transactions': num_transactions, |
| 'period_length': round(period_length, 2), |
| 'period_unit': period_name, |
| 'transaction_velocity': round(velocity, 2), |
| 'first_transaction': account_df['TransactionDate'].min().strftime('%Y-%m-%d'), |
| 'last_transaction': account_df['TransactionDate'].max().strftime('%Y-%m-%d'), |
| 'total_amount': account_df['AmountINR'].sum() |
| }) |
| |
| metrics_df = pd.DataFrame(account_metrics) |
| threshold = metrics_df['transaction_velocity'].quantile(1 - (top_percent/100)) |
| high_velocity_accounts = metrics_df[metrics_df['transaction_velocity'] > threshold].to_dict('records') |
| |
| high_velocity_accounts.sort(key=lambda x: x['transaction_velocity'], reverse=True) |
| |
| |
| try: |
| conn = psycopg2.connect( |
| host="13.126.242.31", |
| database="aml", |
| user="dev_cbs_admin", |
| password="Finovate@2023" |
| ) |
| cur = conn.cursor() |
| |
| |
| rows_inserted = 0 |
| |
| |
| for account in high_velocity_accounts: |
| insert_query = """ |
| INSERT INTO velocity ( |
| account_number, beneficiary_name, total_transactions, |
| period_length, period_unit, transaction_velocity, |
| first_transaction, last_transaction, total_amount |
| ) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s) |
| """ |
| |
| values = ( |
| account['account_number'], |
| account['beneficiary_name'], |
| account['total_transactions'], |
| account['period_length'], |
| account['period_unit'], |
| account['transaction_velocity'], |
| account['first_transaction'], |
| account['last_transaction'], |
| account['total_amount'] |
| ) |
| |
| cur.execute(insert_query, values) |
| rows_inserted += 1 |
| |
| |
| conn.commit() |
| print(f"Response code: 200") |
| |
| |
| sigma_response = await sigma() |
| |
| return { |
| "status": "success", |
| "code": 200, |
| "velocity_analysis_complete": True, |
| "rows_inserted_to_velocity": rows_inserted, |
| "sigma_processing": { |
| "status": sigma_response.get("status"), |
| "total_matching_rows": sigma_response.get("total_matching_rows", 0), |
| "rows_inserted_to_backup": sigma_response.get("rows_inserted_to_backup", 0) |
| } |
| } |
| |
| except Exception as e: |
| print(f"Database error: {str(e)}") |
| print(f"Response code: 400") |
| return { |
| "status": "error", |
| "code": 400, |
| "message": f"Failed to store velocity analysis: {str(e)}" |
| } |
| finally: |
| if 'cur' in locals(): |
| cur.close() |
| if 'conn' in locals(): |
| conn.close() |
|
|
| @app.post("/risk_analysis") |
| async def risk_analysis( |
| file: UploadFile = File(...), |
| daily_threshold: int = Query(10, description="Daily transaction threshold"), |
| weekly_threshold: int = Query(100, description="Weekly transaction threshold"), |
| monthly_threshold: int = Query(1000, description="Monthly transaction threshold") |
| ): |
| print("Starting risk-weighted analysis...") |
| |
| |
| file_content = await file.read() |
| excel_io = io.BytesIO(file_content) |
| df = pd.read_excel(excel_io) |
| |
| df['TransactionDate'] = pd.to_datetime(df['TransactionDate']) |
| |
| risk_flags = { |
| "abnormal_amounts": [], |
| "suspicious_volume": [], |
| "dormant_to_active": [] |
| } |
| |
| |
| def analyze_abnormal_amounts(transactions): |
| mean_amount = transactions['AmountINR'].mean() |
| std_amount = transactions['AmountINR'].std() |
| |
| if std_amount == 0: |
| return [] |
| |
| transactions['z_score'] = (transactions['AmountINR'] - mean_amount) / std_amount |
| suspicious = transactions[abs(transactions['z_score']) > 3].copy() |
| |
| results = [] |
| for _, txn in suspicious.iterrows(): |
| results.append({ |
| 'transaction_date': txn['TransactionDate'].strftime('%Y-%m-%d %H:%M:%S'), |
| 'account_number': txn['BeneficiaryAccountNumber'], |
| 'amount': txn['AmountINR'], |
| 'z_score': round(txn['z_score'], 2), |
| 'transaction_details': {k: (None if pd.isna(v) else v) for k, v in txn.to_dict().items()} |
| }) |
| return results |
| |
| |
| def analyze_suspicious_volume(transactions, period='D'): |
| if period == 'D': |
| grouped = transactions.groupby([transactions['TransactionDate'].dt.date, 'BeneficiaryAccountNumber']).size() |
| threshold = daily_threshold |
| elif period == 'W': |
| grouped = transactions.groupby([transactions['TransactionDate'].dt.isocalendar().week, 'BeneficiaryAccountNumber']).size() |
| threshold = weekly_threshold |
| else: |
| grouped = transactions.groupby([transactions['TransactionDate'].dt.to_period('M'), 'BeneficiaryAccountNumber']).size() |
| threshold = monthly_threshold |
| |
| suspicious = [] |
| for account in transactions['BeneficiaryAccountNumber'].unique(): |
| |
| if (len(grouped) == 0) or (account not in grouped.index.levels[1]): |
| continue |
| account_txns = grouped.xs(account, level=1) |
| if len(account_txns) > 0: |
| p80 = account_txns.quantile(0.80) |
| p99 = account_txns.quantile(0.99) |
| if p99 - p80 < threshold: |
| suspicious.append({ |
| 'account_number': account, |
| 'period': period, |
| 'p80_transactions': int(p80), |
| 'p99_transactions': int(p99), |
| 'difference': int(p99 - p80) |
| }) |
| return suspicious |
| |
| |
| def analyze_dormant_to_active(transactions): |
| account_velocities = [] |
| for account in transactions['BeneficiaryAccountNumber'].unique(): |
| account_txns = transactions[transactions['BeneficiaryAccountNumber'] == account] |
| date_range = (account_txns['TransactionDate'].max() - account_txns['TransactionDate'].min()).days |
| num_transactions = len(account_txns) |
| velocity = num_transactions / (date_range if date_range > 0 else 1) |
| account_velocities.append({ |
| 'account_number': account, |
| 'velocity': velocity |
| }) |
| |
| if not account_velocities: |
| return [] |
| |
| velocities_df = pd.DataFrame(account_velocities) |
| mean_velocity = velocities_df['velocity'].mean() |
| |
| suspicious = [] |
| for acc in account_velocities: |
| if acc['velocity'] > mean_velocity: |
| account_txns = transactions[transactions['BeneficiaryAccountNumber'] == acc['account_number']] |
| suspicious.append({ |
| 'account_number': acc['account_number'], |
| 'velocity': round(acc['velocity'], 2), |
| 'mean_velocity': round(mean_velocity, 2), |
| 'total_transactions': len(account_txns), |
| 'first_transaction': account_txns['TransactionDate'].min().strftime('%Y-%m-%d'), |
| 'last_transaction': account_txns['TransactionDate'].max().strftime('%Y-%m-%d') |
| }) |
| return suspicious |
| |
| |
| risk_flags["abnormal_amounts"] = analyze_abnormal_amounts(df) |
| |
| |
| risk_flags["suspicious_volume"].extend(analyze_suspicious_volume(df, 'D')) |
| risk_flags["suspicious_volume"].extend(analyze_suspicious_volume(df, 'W')) |
| risk_flags["suspicious_volume"].extend(analyze_suspicious_volume(df, 'M')) |
| |
| risk_flags["dormant_to_active"] = analyze_dormant_to_active(df) |
| |
| return { |
| "total_transactions": len(df), |
| "analysis_date": datetime.now().strftime('%Y-%m-%d %H:%M:%S'), |
| "thresholds": { |
| "daily": daily_threshold, |
| "weekly": weekly_threshold, |
| "monthly": monthly_threshold, |
| "z_score": 3 |
| }, |
| "risk_flags": risk_flags |
| } |
|
|
| @app.get("/sanctions_matches", description="Get all sanctions matches from database") |
| async def get_sanctions_matches(): |
| try: |
| conn = psycopg2.connect( |
| host="13.126.242.31", |
| database="aml", |
| user="dev_cbs_admin", |
| password="Finovate@2023" |
| ) |
| cur = conn.cursor(cursor_factory=RealDictCursor) |
| |
| |
| select_query = """ |
| SELECT * FROM sanctions_matches |
| ORDER BY created_at DESC |
| """ |
| cur.execute(select_query) |
| matches = cur.fetchall() |
| |
| return { |
| "status": "success", |
| "table": "sanctions_matches", |
| "total_matches": len(matches), |
| "matches": matches |
| } |
| |
| except Exception as e: |
| print(f"Database error: {str(e)}") |
| return { |
| "status": "error", |
| "message": f"Failed to fetch sanctions matches: {str(e)}" |
| } |
| finally: |
| if 'cur' in locals(): |
| cur.close() |
| if 'conn' in locals(): |
| conn.close() |
|
|
| @app.get("/large_transactions_list", description="Get all large transactions from database") |
| async def get_large_transactions(): |
| try: |
| conn = psycopg2.connect( |
| host="13.126.242.31", |
| database="aml", |
| user="dev_cbs_admin", |
| password="Finovate@2023" |
| ) |
| cur = conn.cursor(cursor_factory=RealDictCursor) |
| |
| |
| select_query = """ |
| SELECT * FROM large_transactions |
| ORDER BY transactiondate DESC |
| """ |
| cur.execute(select_query) |
| transactions = cur.fetchall() |
| |
| return { |
| "status": "success", |
| "table": "large_transactions", |
| "total_transactions": len(transactions), |
| "transactions": transactions |
| } |
| |
| except Exception as e: |
| print(f"Database error: {str(e)}") |
| return { |
| "status": "error", |
| "message": f"Failed to fetch large transactions: {str(e)}" |
| } |
| finally: |
| if 'cur' in locals(): |
| cur.close() |
| if 'conn' in locals(): |
| conn.close() |
|
|
| @app.post("/init_db") |
| async def init_db(): |
| try: |
| conn = psycopg2.connect( |
| host="13.126.242.31", |
| database="aml", |
| user="dev_cbs_admin", |
| password="Finovate@2023" |
| ) |
| cur = conn.cursor() |
| |
| |
| cur.execute(""" |
| DROP TABLE IF EXISTS velocity; |
| CREATE TABLE velocity ( |
| id SERIAL PRIMARY KEY, |
| account_number VARCHAR(100), |
| beneficiary_name TEXT, |
| total_transactions INTEGER, |
| period_length DECIMAL(10,2), |
| period_unit VARCHAR(20), |
| transaction_velocity DECIMAL(10,2), |
| first_transaction DATE, |
| last_transaction DATE, |
| total_amount DECIMAL(15,2), |
| created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, |
| status VARCHAR(20) DEFAULT 'pending' |
| ); |
| """) |
| |
| conn.commit() |
| return {"status": "success", "message": "Database initialized successfully"} |
| |
| except Exception as e: |
| print(f"Database error: {str(e)}") |
| return {"status": "error", "message": str(e)} |
| finally: |
| if 'cur' in locals(): |
| cur.close() |
| if 'conn' in locals(): |
| conn.close() |
|
|
| from pydantic import BaseModel |
|
|
| class StatusUpdate(BaseModel): |
| beneficiary_name: str |
| new_status: str |
| table_name: str |
|
|
| @app.post("/update_status") |
| async def update_status(update: StatusUpdate): |
| valid_tables = ['sanctions_matches', 'large_transactions', 'velocity'] |
| |
| if update.table_name not in valid_tables: |
| return { |
| "status": "error", |
| "message": f"Invalid table name. Must be one of: {', '.join(valid_tables)}" |
| } |
| |
| try: |
| conn = psycopg2.connect( |
| host="13.126.242.31", |
| database="aml", |
| user="dev_cbs_admin", |
| password="Finovate@2023" |
| ) |
| cur = conn.cursor() |
| |
| |
| if update.table_name == 'sanctions_matches': |
| update_query = """ |
| UPDATE sanctions_matches |
| SET status = %s |
| WHERE matched_name = %s |
| """ |
| else: |
| beneficiary_column = ( |
| "beneficiaryname" if update.table_name == 'large_transactions' |
| else "beneficiary_name" if update.table_name == 'velocity' |
| else "beneficiary_name" |
| ) |
| |
| |
| update_query = f""" |
| UPDATE {update.table_name} |
| SET status = %s |
| WHERE {beneficiary_column} = %s |
| """ |
| |
| cur.execute(update_query, (update.new_status, update.beneficiary_name)) |
| rows_updated = cur.rowcount |
| |
| conn.commit() |
| |
| if rows_updated > 0: |
| return { |
| "status": "success", |
| "message": f"Updated {rows_updated} records in {update.table_name}" |
| } |
| else: |
| return { |
| "status": "warning", |
| "message": f"No records found with beneficiary name '{update.beneficiary_name}' in {update.table_name}" |
| } |
| |
| except Exception as e: |
| print(f"Database error: {str(e)}") |
| return { |
| "status": "error", |
| "message": f"Failed to update status: {str(e)}" |
| } |
| finally: |
| if 'cur' in locals(): |
| cur.close() |
| if 'conn' in locals(): |
| conn.close() |
|
|
| @app.get("/load_data_based_on_status/{table_name}/{status}") |
| async def load_data_based_on_status( |
| table_name: str, |
| status: str |
| ): |
| valid_tables = ['sanctions_matches', 'large_transactions', 'velocity'] |
| |
| if table_name not in valid_tables: |
| return { |
| "status": "error", |
| "message": f"Invalid table name. Must be one of: {', '.join(valid_tables)}" |
| } |
| |
| try: |
| conn = psycopg2.connect( |
| host="13.126.242.31", |
| database="aml", |
| user="dev_cbs_admin", |
| password="Finovate@2023" |
| ) |
| cur = conn.cursor(cursor_factory=RealDictCursor) |
| |
| |
| |
| order_by_clause = ( |
| "ORDER BY transactiondate DESC" if table_name == 'large_transactions' |
| else "ORDER BY created_at DESC" if table_name == 'sanctions_matches' |
| else "" |
| ) |
| |
| select_query = f""" |
| SELECT * FROM {table_name} |
| WHERE status = %s |
| {order_by_clause} |
| """ |
| |
| cur.execute(select_query, (status,)) |
| records = cur.fetchall() |
| |
| return { |
| "status": "success", |
| "table": table_name, |
| "filtered_status": status, |
| "total_records": len(records), |
| "records": records |
| } |
| |
| except Exception as e: |
| print(f"Database error: {str(e)}") |
| return { |
| "status": "error", |
| "message": f"Failed to fetch records: {str(e)}" |
| } |
| finally: |
| if 'cur' in locals(): |
| cur.close() |
| if 'conn' in locals(): |
| conn.close() |
|
|
| @app.get("/suspicious_velocity", description="Get all suspicious velocity records from database") |
| async def get_suspicious_velocity(): |
| try: |
| conn = psycopg2.connect( |
| host="13.126.242.31", |
| database="aml", |
| user="dev_cbs_admin", |
| password="Finovate@2023" |
| ) |
| cur = conn.cursor(cursor_factory=RealDictCursor) |
| |
| |
| select_query = """ |
| SELECT * FROM velocity |
| ORDER BY transaction_velocity DESC |
| """ |
| cur.execute(select_query) |
| velocity_records = cur.fetchall() |
| |
| return { |
| "status": "success", |
| "table": "velocity", |
| "total_records": len(velocity_records), |
| "records": velocity_records |
| } |
| |
| except Exception as e: |
| print(f"Database error: {str(e)}") |
| return { |
| "status": "error", |
| "message": f"Failed to fetch velocity records: {str(e)}" |
| } |
| finally: |
| if 'cur' in locals(): |
| cur.close() |
| if 'conn' in locals(): |
| conn.close() |
|
|
| @app.get("/sigma") |
| async def sigma(): |
| try: |
| |
| conn = psycopg2.connect( |
| host="13.126.242.31", |
| database="aml", |
| user="dev_cbs_admin", |
| password="Finovate@2023" |
| ) |
| cur = conn.cursor(cursor_factory=RealDictCursor) |
| |
| |
| select_query = """ |
| SELECT DISTINCT account_number, beneficiary_name |
| FROM velocity |
| ORDER BY account_number |
| """ |
| cur.execute(select_query) |
| velocity_data = cur.fetchall() |
| velocity_accounts = [row['account_number'] for row in velocity_data] |
| velocity_details = [{'account_number': row['account_number'], 'beneficiary_name': row['beneficiary_name']} for row in velocity_data] |
| |
| if not velocity_accounts: |
| return { |
| "status": "warning", |
| "message": "No account numbers found in velocity table", |
| "accounts": [], |
| "velocity_details": [] |
| } |
| |
| |
| try: |
| df = pd.read_excel('temp1.xlsx') |
| |
| |
| filtered_df = pd.DataFrame() |
| |
| |
| for detail in velocity_details: |
| matches = df[ |
| (df['BeneficiaryAccountNumber'] == detail['account_number']) & |
| (df['BeneficiaryName'] == detail['beneficiary_name']) |
| ] |
| filtered_df = pd.concat([filtered_df, matches]) |
| |
| if len(filtered_df) == 0: |
| return { |
| "status": "warning", |
| "message": "No matching transactions found in temp1.xlsx", |
| "accounts_from_velocity": velocity_accounts, |
| "velocity_details": velocity_details, |
| "excel_headers": list(df.columns), |
| "matching_rows": [] |
| } |
| |
| |
| cur.execute(""" |
| SELECT column_name |
| FROM information_schema.columns |
| WHERE table_name = 'backup_data' |
| ORDER BY ordinal_position |
| """) |
| backup_data_columns = [row['column_name'] for row in cur.fetchall()] |
| print("Database columns:", backup_data_columns) |
|
|
| |
| cur.execute(""" |
| SELECT column_name |
| FROM information_schema.columns |
| WHERE table_name = 'backup_data' |
| ORDER BY ordinal_position |
| """) |
| backup_data_columns = [row['column_name'] for row in cur.fetchall()] |
| print("Database columns:", backup_data_columns) |
|
|
| |
| cur.execute(""" |
| SELECT column_name |
| FROM information_schema.columns |
| WHERE table_name = 'backup_data' |
| ORDER BY ordinal_position |
| """) |
| backup_data_columns = [row['column_name'] for row in cur.fetchall()] |
| print("Database columns:", backup_data_columns) |
|
|
| |
| column_mapping = { |
| 'TransactionDate': 'transaction_date', |
| 'TransactionID': 'transaction_id', |
| 'Customer Name': 'customer_name', |
| 'CustomerId': 'customer_id', |
| 'BeneficiaryAccountNumber': 'beneficiary_account_number', |
| 'BeneficiaryName': 'beneficiary_name', |
| 'AmountINR': 'amount_inr', |
| 'SenderName': 'sender_name', |
| 'SenderAccountNumber': 'sender_account_number', |
| 'TransactionType': 'transaction_type', |
| 'RelationshipFlag': 'relationship_flag', |
| 'TransactionTime': 'transaction_time', |
| 'TransactionStatus': 'transaction_status', |
| 'SenderMobileNumber': 'sender_mobile_number', |
| 'SenderIFSC': 'sender_ifsc', |
| 'BeneficiaryMobileNumber': 'beneficiary_mobile_number', |
| 'BeneficiaryIFSC': 'beneficiary_ifsc', |
| 'InstrumentType': 'instrument_type', |
| 'InstrumentIDNumber': 'instrument_id_number', |
| 'InstrumentIssuerInstituteName': 'instrument_issuer_institute_name', |
| 'Purpose': 'purpose', |
| 'TransactionType1': 'transaction_type_1', |
| 'TransactionAmount': 'transaction_amount', |
| 'AmountFC': 'amount_fc', |
| 'FCCode': 'fc_code', |
| 'Narration': 'narration', |
| 'SenderVPA': 'sender_vpa', |
| 'DeclarationStatingSenderVPAIsNotAvailable': 'declaration_stating_sender_vpa_is_not_available', |
| 'BeneficiaryVPA': 'beneficiary_vpa', |
| 'MerchantCategoryCode': 'merchant_category_code', |
| 'BeneficiaryAccountType': 'beneficiary_account_type', |
| 'CardNumber': 'card_number', |
| 'MaskedOrNonMasked': 'masked_or_non_masked', |
| 'MerchantName': 'merchant_name', |
| 'MerchantID': 'merchant_id', |
| 'MerchantPoSOrATMID': 'merchant_pos_or_atm_id', |
| 'MerchantCountry': 'merchant_country', |
| 'AddressLine1': 'address_line1', |
| 'Locality': 'locality', |
| 'Country': 'country', |
| 'PINCode': 'pin_code', |
| 'State': 'state', |
| 'District': 'district', |
| 'CityVillageTown': 'city_village_town', |
| 'AcquiringBank': 'acquiring_bank', |
| 'IssuingCountry': 'issuing_country', |
| 'IssuingBank': 'issuing_bank', |
| 'TransactionCurrencyCode': 'transaction_currency_code', |
| 'CustomerName': 'customer_name', |
| 'AccountStatus': 'account_status' |
| } |
|
|
| |
| rows_inserted = 0 |
| for _, row in filtered_df.iterrows(): |
| |
| db_values = {} |
| for excel_col, db_col in column_mapping.items(): |
| if excel_col in row: |
| db_values[db_col] = row[excel_col] |
|
|
| |
| columns = list(db_values.keys()) |
| placeholders = ['%s'] * len(columns) |
| |
| insert_query = f""" |
| INSERT INTO backup_data ( |
| {', '.join(columns)} |
| ) VALUES ( |
| {', '.join(placeholders)} |
| ) |
| """ |
| |
| |
| values = [db_values[col] for col in columns] |
| |
| cur.execute(insert_query, values) |
| rows_inserted += 1 |
| |
| conn.commit() |
| |
| |
| matching_rows = [] |
| for _, row in filtered_df.iterrows(): |
| row_dict = {} |
| for col, val in row.items(): |
| |
| if pd.isna(val): |
| row_dict[col] = None |
| elif isinstance(val, float) and (np.isinf(val) or np.isneginf(val)): |
| row_dict[col] = str(val) |
| else: |
| row_dict[col] = val |
| matching_rows.append(row_dict) |
|
|
| return { |
| "status": "success", |
| "total_accounts_in_velocity": len(velocity_accounts), |
| "accounts_from_velocity": velocity_accounts, |
| "velocity_details": velocity_details, |
| "excel_headers": list(df.columns), |
| "total_matching_rows": len(filtered_df), |
| "matching_rows": matching_rows, |
| "rows_inserted_to_backup": rows_inserted |
| } |
| |
| except FileNotFoundError: |
| return { |
| "status": "error", |
| "message": "temp1.xlsx not found. Please upload a file first using velocity_analysis or large_transactions API", |
| "accounts_from_velocity": velocity_accounts |
| } |
| |
| except Exception as e: |
| print(f"Database error: {str(e)}") |
| return { |
| "status": "error", |
| "message": f"Failed to fetch account numbers: {str(e)}" |
| } |
| finally: |
| if 'cur' in locals(): |
| cur.close() |
| if 'conn' in locals(): |
| conn.close() |
|
|
| @app.get("/") |
| async def root(): |
| return { |
| "message": "Welcome to the Search API", |
| "endpoints": [ |
| "/search/individuals", |
| "/search/entities", |
| "/query", |
| "/large_transactions", |
| "/large_transactions_list", |
| "/sanctions_matches", |
| "/sanction_list_2", |
| "/velocity_analysis", |
| "/risk_analysis" |
| ] |
| } |
|
|
| if __name__ == "__main__": |
| port = int(os.environ.get("PORT", 7860)) |
| uvicorn.run(app, host="0.0.0.0", port=port) |