first commit
This commit is contained in:
77
local/worker.py
Normal file
77
local/worker.py
Normal file
@@ -0,0 +1,77 @@
|
||||
"""Transactional outbox worker. Fake receipts persist; no messages leave the stack."""
|
||||
import json
|
||||
import logging
|
||||
import threading
|
||||
import time
|
||||
from http.server import BaseHTTPRequestHandler, HTTPServer
|
||||
from psycopg.types.json import Jsonb
|
||||
from .adapters import FakeTiny, FakeWhatsApp, LocalS3Storage, require_local
|
||||
from .db import connect
|
||||
from .scanning import ClamAV, scan_loop
|
||||
|
||||
require_local()
|
||||
adapters = {'tiny': FakeTiny(), 'whatsapp': FakeWhatsApp()}
|
||||
last_tick = 0.0
|
||||
last_cleanup = 0.0
|
||||
storage = LocalS3Storage()
|
||||
scan_thread = None
|
||||
|
||||
def cleanup():
|
||||
"""Delete only expired object bytes; retain order/file metadata and history."""
|
||||
with connect() as c:
|
||||
rows = c.execute("""SELECT * FROM dtf_local.uploads WHERE purged_at IS NULL
|
||||
AND (expires_at<=now() OR (NOT complete AND created_at<now()-interval '1 day'))
|
||||
ORDER BY created_at FOR UPDATE SKIP LOCKED LIMIT 50""").fetchall()
|
||||
for row in rows:
|
||||
storage.discard(row['object_key'],row['multipart_id'],row['complete'])
|
||||
c.execute('UPDATE dtf_local.uploads SET purged_at=now() WHERE id=%s', (row['id'],))
|
||||
c.execute('DELETE FROM dtf_local.sessions WHERE expires_at<=now()')
|
||||
c.execute('DELETE FROM dtf_local.operator_sessions WHERE expires_at<=now()')
|
||||
c.execute("DELETE FROM dtf_local.login_attempts WHERE started_at<now()-interval '1 day'")
|
||||
c.execute("DELETE FROM dtf_local.security_events WHERE created_at<now()-interval '30 days'")
|
||||
|
||||
def tick():
|
||||
global last_tick
|
||||
with connect() as c:
|
||||
job = c.execute('SELECT * FROM dtf_local.outbox WHERE delivered_at IS NULL AND available_at <= now() ORDER BY id FOR UPDATE SKIP LOCKED LIMIT 1').fetchone()
|
||||
if job:
|
||||
try:
|
||||
receipt = adapters[job['provider']].deliver(job['event_key'], job['payload'])
|
||||
c.execute('UPDATE dtf_local.outbox SET delivered_at=now(), receipt=%s, attempts=attempts+1, last_error=NULL WHERE id=%s', (Jsonb(receipt),job['id']))
|
||||
except Exception as exc:
|
||||
delay = min(1800, 2**min(job['attempts']+1, 10))
|
||||
c.execute("UPDATE dtf_local.outbox SET attempts=attempts+1, last_error=%s, available_at=now() + %s * interval '1 second' WHERE id=%s", (str(exc),delay,job['id']))
|
||||
last_tick = time.monotonic()
|
||||
|
||||
def loop():
|
||||
global last_cleanup
|
||||
while True:
|
||||
try:
|
||||
tick()
|
||||
if time.monotonic()-last_cleanup > 60:
|
||||
cleanup()
|
||||
last_cleanup = time.monotonic()
|
||||
except Exception:
|
||||
logging.exception('Local worker tick failed')
|
||||
time.sleep(1)
|
||||
|
||||
class Health(BaseHTTPRequestHandler):
|
||||
def do_GET(self):
|
||||
scanner = False
|
||||
try:
|
||||
scanner = bool(scan_thread and scan_thread.is_alive() and ClamAV().ping())
|
||||
except Exception:
|
||||
pass
|
||||
healthy = time.monotonic()-last_tick < 15 and scanner
|
||||
self.send_response(200 if self.path == '/health' and healthy else 503)
|
||||
self.end_headers()
|
||||
self.wfile.write(json.dumps({'worker': 'ok' if healthy else 'unavailable',
|
||||
'scanner': 'ok' if scanner else 'unavailable'}).encode())
|
||||
def log_message(self, *args):
|
||||
pass
|
||||
|
||||
if __name__ == '__main__':
|
||||
scan_thread = threading.Thread(target=scan_loop,args=(storage,),daemon=True)
|
||||
scan_thread.start()
|
||||
threading.Thread(target=loop, daemon=True).start()
|
||||
HTTPServer(('0.0.0.0',8002),Health).serve_forever()
|
||||
Reference in New Issue
Block a user