All checks were successful
Build and deploy / Validate source (push) Successful in 9s
Build and deploy / Integration suite on a real stack (push) Successful in 2m49s
Build and deploy / Secret scan and release gate (push) Successful in 9s
Build and deploy / Publish images and notify Portainer (push) Has been skipped
Tiny v3 replaces the v2 token adapter. An operator connects Tiny once from the Kanban; the callback is authorised by a single-use state, because Tiny's cross-site redirect does not carry the SameSite=Strict operator cookie. Tokens are kept in provider_tokens, the refresh token rotates under a row lock, and the worker keeps the connection alive while order creation is off. Orders find or create the customer's contact by CNPJ, then POST /pedidos with product ids from TINY_PRODUCT_TEXTIL_FOLHA, _TEXTIL_AVULSA, _UV_FOLHA and _UV_AVULSA and numeroOrdemCompra DTF-<number>; a retry searches the customer's recent orders for that number first. The product settings avoid a _FILE suffix, which the secrets loader reads as a secret file path. Production passes the application credentials through but keeps TINY_ADAPTER fake: Tiny has no sandbox, so creating real orders waits for a supervised test. compose.providers.yaml gives the local API and worker an internet route for provider testing; the default local stack still has none. Verified with the full CI integration sequence locally, including the new tiny_oauth_test against the real database. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
114 lines
4.8 KiB
Python
114 lines
4.8 KiB
Python
"""Transactional outbox worker. Fake receipts persist; no messages leave the stack."""
|
|
import json
|
|
import logging
|
|
import os
|
|
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 .printjobs import render_loop
|
|
from .scanning import ClamAV, scan_loop
|
|
|
|
load_secret_files()
|
|
require_runtime()
|
|
adapters = {'tiny': FakeTiny(), 'whatsapp': FakeWhatsApp()}
|
|
# Not yet confirmed against the client's account; see app/tiny.py.
|
|
if os.environ.get('TINY_ADAPTER') == 'tiny':
|
|
from .tiny import TinyOrders
|
|
adapters['tiny'] = TinyOrders()
|
|
last_tick = 0.0
|
|
last_cleanup = 0.0
|
|
last_tiny_keepalive = 0.0
|
|
storage = LocalS3Storage()
|
|
scan_thread = None
|
|
render_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 tiny_keepalive():
|
|
"""Keep a connected Tiny authorised even while no orders are sent.
|
|
|
|
The refresh token expires unless it is used; access_token() refreshes (and
|
|
rotates) only when the access token is about to expire, so asking every few
|
|
minutes renews the connection roughly once per access-token lifetime.
|
|
"""
|
|
global last_tiny_keepalive
|
|
if time.monotonic()-last_tiny_keepalive < 600:
|
|
return
|
|
last_tiny_keepalive = time.monotonic()
|
|
from . import tiny
|
|
if not tiny.configured():
|
|
return
|
|
try:
|
|
tiny.TinyAuth().access_token()
|
|
except tiny.TinyNotConnected:
|
|
pass
|
|
except Exception:
|
|
logging.exception('Tiny connection refresh failed')
|
|
|
|
def loop():
|
|
global last_cleanup
|
|
while True:
|
|
try:
|
|
tick()
|
|
if time.monotonic()-last_cleanup > 60:
|
|
cleanup()
|
|
last_cleanup = time.monotonic()
|
|
tiny_keepalive()
|
|
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
|
|
renderer = bool(render_thread and render_thread.is_alive())
|
|
healthy = time.monotonic()-last_tick < 15 and scanner and renderer
|
|
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',
|
|
'print_files': 'ok' if renderer 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()
|
|
render_thread = threading.Thread(target=render_loop,args=(storage,),daemon=True)
|
|
render_thread.start()
|
|
threading.Thread(target=loop, daemon=True).start()
|
|
HTTPServer(('0.0.0.0',8002),Health).serve_forever()
|