-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathAdvance_settlement_script
More file actions
278 lines (250 loc) · 11.4 KB
/
Copy pathAdvance_settlement_script
File metadata and controls
278 lines (250 loc) · 11.4 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
from datetime import datetime
from decimal import Decimal
import psycopg2
# === DATABASE CONNECTIONS ===
conn_main = psycopg2.connect(
dbname='postgres',
user='postgres',
password='postgres',
host='localhost',
port='5432'
)
cur_main = conn_main.cursor()
conn_local = psycopg2.connect(
dbname='postgres',
user='postgres',
password='postgres',
host='localhost',
port='5432'
)
cur_local = conn_local.cursor()
# === LOGGING FUNCTIONS ===
def log_full_sql_query_v2(local_cursor, consumer_code, tenant_id, receipt_number, advance_amount,
demand_id, demand_detail_id, query_valid_advances,
check_future_payment_query, select_advance_demand_query,
select_future_demand_query, cancel_queries, update_advance_query,
update_demand_status_query, expire_bill_query, update_bill_queries):
cancel_queries_str = "\n".join(cancel_queries)
update_bill_queries_str = "\n".join(update_bill_queries)
local_cursor.execute("""
INSERT INTO full_sql_query_log (
consumer_code, tenant_id, receipt_number, advance_amount,
demand_id, demand_detail_id, query_valid_advances,
check_future_payment_query, select_advance_demand_query,
select_future_demand_query, cancel_queries,
update_advance_query, update_demand_status_query,
expire_bill_query, update_bill_queries, failure_reason
) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, NULL);
""", (
consumer_code, tenant_id, receipt_number, advance_amount,
demand_id, demand_detail_id, query_valid_advances.strip(),
check_future_payment_query.strip(), select_advance_demand_query.strip(),
select_future_demand_query.strip(), cancel_queries_str,
update_advance_query.strip(), update_demand_status_query.strip(),
expire_bill_query.strip(), update_bill_queries_str
))
def log_failure_reason_only(local_cursor, consumer_code, tenant_id, receipt_number, advance_amount, failure_reason):
local_cursor.execute("""
INSERT INTO full_sql_query_log (
consumer_code, tenant_id, receipt_number, advance_amount, failure_reason
) VALUES (%s, %s, %s, %s, %s);
""", (
consumer_code, tenant_id, receipt_number, advance_amount, failure_reason
))
# === BACKUP FUNCTION ===
def backup_before_after_values(main_cursor, local_cursor, demand_id, demand_detail_id, bill_ids, advance_amount):
# Backup egbs_demand_v1
main_cursor.execute(f"SELECT ispaymentcompleted FROM egbs_demand_v1 WHERE id = %s", (demand_id,))
before_ispaymentcompleted = main_cursor.fetchone()[0]
# Backup egbs_demanddetail_v1
main_cursor.execute(f"""
SELECT taxheadcode, taxamount, collectionamount
FROM egbs_demanddetail_v1
WHERE id = %s
""", (demand_detail_id,))
taxhead_code, before_taxamount, before_collectionamount = main_cursor.fetchone()
# Backup egbs_bill_v1 (first active bill only)
before_bill_status = after_bill_status = None
bill_id = None
if bill_ids:
bill_id = bill_ids[0]
main_cursor.execute("SELECT status FROM egbs_bill_v1 WHERE id = %s", (bill_id,))
before_bill_status = main_cursor.fetchone()[0]
# === AFTER VALUES (after update) ===
after_ispaymentcompleted = True
after_taxamount = Decimal('0.00')
after_collectionamount = Decimal('0.00')
if bill_id:
after_bill_status = "EXPIRED"
# INSERT BACKUP ROW
local_cursor.execute("""
INSERT INTO advance_update_history (
consumer_code, tenant_id, demand_id,
before_ispaymentcompleted, after_ispaymentcompleted,
demand_detail_id, taxhead_code,
before_taxamount, after_taxamount,
before_collectionamount, after_collectionamount,
bill_id, before_bill_status, after_bill_status,
updated_at
) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, NOW());
""", (
consumer_code, tenant_id, demand_id,
before_ispaymentcompleted, after_ispaymentcompleted,
demand_detail_id, taxhead_code,
before_taxamount, after_taxamount,
before_collectionamount, after_collectionamount,
bill_id, before_bill_status, after_bill_status
))
# === PROCESSING SCRIPT ===
cur_local.execute("""
SELECT id, consumer_code, tenant_id, business_service
FROM advance_payment_details
WHERE processed = FALSE
ORDER BY id;
""")
queue_items = cur_local.fetchall()
for queue_id, consumer_code, tenant_id, business_service in queue_items:
try:
query_valid_advances = f"""
SELECT bill.consumercode, pd.receiptnumber, p.totalamountpaid,
p.paymentstatus, bill.status, bill.totalamount,
pd.receiptdate, bill.tenantid, bill.businessservice
FROM egcl_payment p
JOIN egcl_paymentdetail pd ON pd.paymentid = p.id
JOIN egcl_bill bill ON bill.id = pd.billid
WHERE bill.tenantid = '{tenant_id}'
AND bill.businessservice = '{business_service}'
AND bill.consumercode = '{consumer_code}'
AND pd.amountpaid > pd.due
AND NOT EXISTS (
SELECT 1 FROM egcl_paymentdetail pd2
JOIN egcl_bill bill2 ON bill2.id = pd2.billid
WHERE bill2.consumercode = bill.consumercode
AND bill2.tenantid = '{tenant_id}'
AND pd2.receiptdate > pd.receiptdate
)
ORDER BY pd.createdtime DESC;
"""
cur_main.execute(query_valid_advances)
row = cur_main.fetchone()
if not row:
raise Exception("No valid advance receipt found.")
consumer_code, receipt_number, total_paid, _, _, total_amount, receipt_date, tenant_id, business_service = row
advance_amount = Decimal(total_paid) - Decimal(total_amount)
if advance_amount <= 0:
raise Exception("Advance amount not positive")
query_check_new_payment = f"""
SELECT * FROM egcl_payment p
JOIN egcl_paymentdetail pd ON pd.paymentid = p.id
JOIN egcl_bill bill ON bill.id = pd.billid
WHERE bill.tenantid = '{tenant_id}'
AND bill.businessservice = '{business_service}'
AND bill.consumercode = '{consumer_code}'
AND pd.receiptdate > '{receipt_date}';
"""
cur_main.execute(query_check_new_payment)
if cur_main.fetchone():
log_failure_reason_only(cur_local, consumer_code, tenant_id, receipt_number, float(advance_amount),
"New payment found after advance")
cur_local.execute("""
UPDATE advance_payment_details
SET processed = TRUE, processed_at = %s, error_reason = %s
WHERE id = %s;
""", (datetime.now(), "New payment found after advance", queue_id))
conn_local.commit()
continue
select_advance_demand_query = f"""
SELECT pd.receiptnumber, dd.demandid, dd.id AS demanddetailsid
FROM egcl_paymentdetail pd
JOIN egcl_bill bill ON bill.id = pd.billid
JOIN egbs_billdetail_v1 bd ON bd.billid = bill.id
JOIN egbs_demanddetail_v1 dd ON bd.demandid = dd.demandid
JOIN egbs_demand_v1 d ON dd.demandid = d.id
WHERE pd.receiptnumber = '{receipt_number}'
AND bd.businessservice = '{business_service}'
AND dd.taxheadcode = '{business_service}_ADVANCE_CARRYFORWARD'
AND TO_CHAR(TO_TIMESTAMP(dd.createdtime / 1000) AT TIME ZONE 'Asia/Kolkata', 'YYYY-MM-DD') =
TO_CHAR(TO_TIMESTAMP(pd.receiptdate / 1000) AT TIME ZONE 'Asia/Kolkata', 'YYYY-MM-DD');
"""
cur_main.execute(select_advance_demand_query)
demand_data = cur_main.fetchone()
if not demand_data:
raise Exception("Advance demand detail not found")
demand_id = demand_data[1]
demand_detail_id = demand_data[2]
# Cancel future demands
cancel_queries = []
select_future_demand_query = f"""
SELECT d.id FROM egbs_demand_v1 d
WHERE d.businessservice = '{business_service}'
AND d.tenantid = '{tenant_id}'
AND d.consumercode = '{consumer_code}'
AND d.status != 'CANCELLED'
AND d.taxperiodfrom > (SELECT taxperiodfrom FROM egbs_demand_v1 WHERE id = '{demand_id}');
"""
cur_main.execute(select_future_demand_query)
for (cancel_id,) in cur_main.fetchall():
cancel_sql = f"UPDATE egbs_demand_v1 SET status = 'CANCELLED' WHERE id = '{cancel_id}';"
cur_main.execute(cancel_sql)
cancel_queries.append(cancel_sql)
# Get bill ID before update
cur_main.execute(f"""
SELECT bill.id
FROM egbs_bill_v1 bill
JOIN egbs_billdetail_v1 bd ON bd.billid = bill.id
WHERE bill.status = 'ACTIVE'
AND bd.businessservice = '{business_service}'
AND bd.consumercode = '{consumer_code}'
AND bill.tenantid = '{tenant_id}';
""")
bill_ids = [row[0] for row in cur_main.fetchall()]
# Backup before-after
backup_before_after_values(cur_main, cur_local, demand_id, demand_detail_id, bill_ids, advance_amount)
# Update advance record
update_advance_query = f"""
UPDATE egbs_demanddetail_v1
SET taxamount = 0.00, collectionamount = 0.00
WHERE id = '{demand_detail_id}';
"""
cur_main.execute(update_advance_query)
update_demand_status_query = f"""
UPDATE egbs_demand_v1
SET ispaymentcompleted = true
WHERE id = '{demand_id}';
"""
cur_main.execute(update_demand_status_query)
# Expire bills
update_bill_queries = []
for bill_id in bill_ids:
sql = f"UPDATE egbs_bill_v1 SET status = 'EXPIRED' WHERE id = '{bill_id}';"
cur_main.execute(sql)
update_bill_queries.append(sql)
conn_main.commit()
# Log SQL
log_full_sql_query_v2(cur_local, consumer_code, tenant_id, receipt_number, float(advance_amount),
demand_id, demand_detail_id, query_valid_advances, query_check_new_payment,
select_advance_demand_query, select_future_demand_query,
cancel_queries, update_advance_query, update_demand_status_query,
select_future_demand_query, update_bill_queries)
cur_local.execute("""
UPDATE advance_payment_details
SET processed = TRUE, processed_at = %s
WHERE id = %s;
""", (datetime.now(), queue_id))
conn_local.commit()
print(f"✅ Processed: {consumer_code} / {receipt_number}")
except Exception as e:
print(f"❌ Error for {consumer_code}: {e}")
conn_main.rollback()
cur_local.execute("""
UPDATE advance_payment_details
SET processed = TRUE, processed_at = %s, error_reason = %s
WHERE id = %s;
""", (datetime.now(), str(e), queue_id))
conn_local.commit()
# === CLEANUP ===
cur_main.close()
conn_main.close()
cur_local.close()
conn_local.close()
print("✅ Queue processing completed.")