refactor: lay the repository out by role
All checks were successful
Build and deploy / Validate source (push) Successful in 7s
Build and deploy / Integration suite on a real stack (push) Successful in 1m25s
Build and deploy / Secret scan and release gate (push) Successful in 5s
Build and deploy / Publish images and notify Portainer (push) Successful in 1m31s
All checks were successful
Build and deploy / Validate source (push) Successful in 7s
Build and deploy / Integration suite on a real stack (push) Successful in 1m25s
Build and deploy / Secret scan and release gate (push) Successful in 5s
Build and deploy / Publish images and notify Portainer (push) Successful in 1m31s
local/ held six unrelated things under a name that stopped being true once it
became the production runtime: the service, the frontend, the tests, the ops
commands, the container definitions and the dependency lock, 65 files with
nothing to tell them apart.
app/ the service: api/ routers, core/ for identity, database, models,
prices and secret loading, and the worker, bootstrap and schema
tests/ the twelve suites, no longer inside the shipped package
ops/ backup, readiness, dependency audit, security summary
infra/ Dockerfiles, gateway templates, ClamAV and storage configuration,
the requirements and their hash lock
web/ the Site, Kanban and portal pages with their scripts
deploy/Dockerfile.api now copies app/ alone, so the tests stop shipping to
production; the local image still carries them, because the suites run inside
the stack's network.
Five kinds of reference had to follow, and each was found by something different
rather than by reading. Imports of the form "from . import db" survived a rewrite
that only matched "from .db import". Tests kept relative imports of modules that
had left the package. A mock.patch target names its module in a string, where no
import rewriting can see it. The browser test resolves a fixture by path. And the
release gate's markers pointed at local/runtime.py and local/worker.py, which is
the decay its new marker test exists to catch — it caught it.
Verified from docker compose down -v: the stack starts, all six integration
suites, both browser suites and the twenty-nine unit tests pass.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
79
app/worker.py
Normal file
79
app/worker.py
Normal file
@@ -0,0 +1,79 @@
|
||||
"""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_runtime
|
||||
from .core.secrets import load as load_secret_files
|
||||
from .core.db import connect
|
||||
from .scanning import ClamAV, scan_loop
|
||||
|
||||
load_secret_files()
|
||||
require_runtime()
|
||||
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