loopable / platform /core /odoo.py
fsanyoto's picture
Deploy AIOS web (React glide grid + FastAPI slice)
0adc1e7 verified
Raw
History Blame Contribute Delete
13.4 kB
"""CFO-OS shared Odoo data layer (READ-ONLY).
Reuses the existing write-guarded client (ffs_dashboard/odoo_client.py) so there is
a single source of truth for the connection. Adds:
- cached singleton client
- read_group / search_read passthroughs with light cleanup
- business constants (teams, excluded partners)
- the "sales universe" base domain (confirmed orders, RI+FFS only)
NOTHING here writes to Odoo — the underlying client hard-blocks writes.
"""
import os
import sys
import queue as _queue
import threading
from concurrent.futures import ThreadPoolExecutor
from contextlib import contextmanager
from functools import lru_cache
from pathlib import Path
# Self-contained read-only client lives at the app root (platform/odoo_client.py) with its
# own .env — so this whole folder can be moved/hosted anywhere without external paths.
_ROOT = Path(__file__).resolve().parents[1] # platform/
if str(_ROOT) not in sys.path:
sys.path.insert(0, str(_ROOT))
from odoo_client import OdooClient # noqa: E402 (read-only, write-guarded)
# ---- business constants (TENANT #0's scope, kept as module attributes for its modules) ----
# These are Royal Imports' config — the default a request gets when nothing binds a tenant.
# The canonical per-tenant home is `harness/tenants.py` Tenant.config; a BOUND tenant (see
# `tenant_scope` below) carries its own team ids / exclusions and never reads these.
TEAM_IDS = [5, 6]
TEAM_NAMES = {5: 'Fisch', 6: 'Royal'}
# Partners to exclude from analytics (not part of RI/FFS). Sourced from env so no
# customer name is ever hardcoded in the (potentially public) repo; set as a Space secret.
EXCLUDE_PARTNER_NAMES = {n.strip() for n in os.environ.get('EXCLUDE_PARTNER_NAMES', '').split(',') if n.strip()}
# Confirmed sale.order.line scope for RI+FFS (the analytics "sales universe").
# ⚠ Documentation of tenant #0's base scope — `sale_line_domain()` below BUILDS its domain from
# `active_team_ids()` (binding-aware) rather than reading this list; edit both or neither.
SALE_LINE_BASE = [
('state', 'in', ['sale', 'done']),
('order_id.team_id', 'in', TEAM_IDS),
('product_id', '!=', False),
]
# Posted customer invoices/credit notes (AR + revenue recognised).
INVOICE_BASE = [
('move_type', 'in', ['out_invoice', 'out_refund']),
('state', '=', 'posted'),
]
@lru_cache(maxsize=1)
def _global_client():
"""TENANT #0's read-only client (env/.env credentials) — the default for every unbound call."""
return OdooClient()
# ---- per-tenant connections (R3 keychain cutover, 2026-08-04) --------------------------------
# One client per tenant slug, keyed by a fingerprint of its credentials so a rotated key
# reconnects on the next call with no restart. The DEFAULT slug maps to the env client above.
# ⛔ FAIL-CLOSED: a bound NON-default tenant with no credentials raises — the env fallback
# would silently answer tenant B's question over tenant A's connection, which is the exact
# failure this section exists to remove.
_DEFAULT_SLUG = 'royal-imports'
_TENANT_CLIENTS = {} # slug -> (creds fingerprint, OdooClient)
_TC_LOCK = threading.Lock()
class TenantOdooUnavailable(RuntimeError):
"""A bound tenant has no Odoo credentials. Callers surface 'connect a source', never
another tenant's data."""
def _client_for(slug, creds):
"""The client for a BOUND tenant. creds None: default slug -> env client; anyone else -> raise."""
slug = str(slug or _DEFAULT_SLUG)
if not creds:
if slug == _DEFAULT_SLUG:
return _global_client()
raise TenantOdooUnavailable(
f"tenant {slug!r} has no Odoo credentials — store them under Settings → Keychains. "
f"(Refusing the environment fallback: that is another tenant's connection.)")
fp = (str(creds.get('url') or ''), str(creds.get('db') or ''),
str(creds.get('user') or ''), str(creds.get('api_key') or ''))
with _TC_LOCK:
cur = _TENANT_CLIENTS.get(slug)
if cur and cur[0] == fp:
return cur[1]
cli = OdooClient(creds=creds) # built OUTSIDE the lock: it authenticates over the wire
with _TC_LOCK:
cur = _TENANT_CLIENTS.get(slug)
if cur and cur[0] == fp:
return cur[1] # another thread won the race — keep ONE client
_TENANT_CLIENTS[slug] = (fp, cli)
return cli
@contextmanager
def tenant_scope(slug, creds=None, team_ids=None, team_names=None, exclude_partner_names=None):
"""Bind THIS THREAD's Odoo context to a tenant: its connection + its scope constants.
Unbound (the default) = tenant #0 via env, which is every existing Streamlit/API caller.
The connector layer (harness/connectors/odoo.py) wraps each query in this, so a second
tenant's keychain credentials reach the wire without any module knowing. Scope constants
are PUSHED here from Tenant.config because core/ may not import harness/ (layering)."""
prev = getattr(_tlocal, 'tenant_ctx', None)
_tlocal.tenant_ctx = {
'slug': str(slug or _DEFAULT_SLUG),
'creds': creds,
'team_ids': list(team_ids) if team_ids else None,
'team_names': dict(team_names) if team_names else None,
'exclude_partner_names': set(exclude_partner_names) if exclude_partner_names else None,
}
try:
yield
finally:
_tlocal.tenant_ctx = prev
def _ctx():
return getattr(_tlocal, 'tenant_ctx', None)
def active_team_ids():
"""The bound tenant's team ids; unbound or default slug -> tenant #0's TEAM_IDS."""
c = _ctx()
if c is not None and c['slug'] != _DEFAULT_SLUG:
return list(c['team_ids'] or [])
return list(TEAM_IDS)
def active_team_names():
c = _ctx()
if c is not None and c['slug'] != _DEFAULT_SLUG:
return dict(c['team_names'] or {})
return dict(TEAM_NAMES)
# ---- concurrent read pool -------------------------------------------------------------------
# Each Odoo XML-RPC call is a round-trip; a heavy drawer fires a dozen of them and they add up to
# many seconds when run one-by-one. `parallel()` runs independent READ tasks across a small pool of
# extra connections, with each task transparently using its own connection via a thread-local — so
# the existing read_group/search_read/sum_field helpers need no change. READ-ONLY only.
_POOL_MAX = 8
_tlocal = threading.local()
_pool_q = _queue.Queue()
_pool_n = 0
_pool_lock = threading.Lock()
def _checkout():
global _pool_n
try:
return _pool_q.get_nowait()
except _queue.Empty:
with _pool_lock:
grow = _pool_n < _POOL_MAX
if grow:
_pool_n += 1
return OdooClient() if grow else _pool_q.get() # else block until one frees
def _checkin(c):
_pool_q.put(c)
def get_odoo():
"""The active read-only client, in priority order:
1. a NON-default tenant binding (tenant_scope) — its own client, ALWAYS. The parallel()
pool below holds tenant-#0 connections, so a pooled client must never answer a bound
tenant's call (that is the cross-tenant leak R3's cutover removes);
2. this task's pooled connection (parallel(), and the app's dedicated-bg-thread pattern);
3. an explicit default-slug binding — the env client;
4. unbound: the shared env singleton. Transparent to every read_group/... caller."""
c = _ctx()
if c is not None and c['slug'] != _DEFAULT_SLUG:
return _client_for(c['slug'], c['creds'])
task_client = getattr(_tlocal, 'client', None)
if task_client is not None:
return task_client
if c is not None:
return _client_for(c['slug'], c['creds'])
return _global_client()
def set_doc_mode(mode):
"""Recognition basis for THIS thread's domain builds: 'order' (all confirmed orders) or 'invoice'
(fully-invoiced orders only). Cached bundles call this before pulling; parallel() propagates it to
each worker. Thread-local, so concurrent user sessions never race on each other's basis."""
_tlocal.doc_mode = 'invoice' if mode == 'invoice' else 'order'
def doc_mode():
return getattr(_tlocal, 'doc_mode', 'order')
def parallel(tasks):
"""Run independent zero-arg READ callables concurrently across the connection pool; returns
results in order. Inside each task, get_odoo() (hence read_group/search_read/sum_field/
search_count) uses that task's own connection. Never use for writes."""
tasks = list(tasks)
if not tasks:
return []
if len(tasks) == 1:
return [tasks[0]()]
caller_mode = doc_mode() # propagate the submitter's basis to each worker
caller_ctx = _ctx() # and the tenant binding — a worker thread must
# answer for the SAME tenant as its submitter
def _run(fn):
c = _checkout()
_tlocal.client = c
_tlocal.doc_mode = caller_mode
_tlocal.tenant_ctx = caller_ctx
try:
return fn()
finally:
_tlocal.client = None
_tlocal.doc_mode = 'order'
_tlocal.tenant_ctx = None
_checkin(c)
with ThreadPoolExecutor(max_workers=min(len(tasks), _POOL_MAX)) as ex:
return list(ex.map(_run, tasks))
def parallel_map(task_dict):
"""parallel() over a {name: callable} mapping -> {name: result}."""
keys = list(task_dict)
return dict(zip(keys, parallel([task_dict[k] for k in keys])))
def m2o_id(v):
return v[0] if isinstance(v, list) and v else None
def m2o_name(v):
return v[1] if isinstance(v, list) and len(v) > 1 else ''
_EXCL_IDS = {} # slug -> tuple of excluded partner ids (resolved once per process)
_EXCL_LOCK = threading.Lock()
def excluded_partner_ids():
"""Resolve the ACTIVE tenant's excluded partner names to ids, for fast domain filtering.
Cached per tenant slug — the old lru_cache(1) would have handed tenant B whatever tenant A
resolved first. Unbound/default = tenant #0's env-sourced names; a bound tenant's names come
from its scope config (none configured -> nothing excluded)."""
c = _ctx()
bound_other = c is not None and c['slug'] != _DEFAULT_SLUG
slug = c['slug'] if bound_other else _DEFAULT_SLUG
names = (c['exclude_partner_names'] or set()) if bound_other else EXCLUDE_PARTNER_NAMES
with _EXCL_LOCK:
if slug in _EXCL_IDS:
return _EXCL_IDS[slug]
if not names:
res = tuple()
else:
rows = get_odoo().search_read('res.partner', [('name', 'in', list(names))], ['id'])
res = tuple(r['id'] for r in rows)
with _EXCL_LOCK:
_EXCL_IDS[slug] = res
return res
def sale_line_domain(date_from=None, date_to=None, team_id=None, extra=None, partner_ids=None):
"""Build a sale.order.line domain in the RI+FFS confirmed-order scope, excluding
the configured house accounts, with optional date window / single team / extra clauses.
partner_ids (a collection, possibly empty) restricts to those customers — the carrier for
the Customer module's Agent filter; None = no partner restriction (an empty set matches none)."""
tids = active_team_ids()
dom = [('state', 'in', ['sale', 'done']), ('product_id', '!=', False)]
if tids:
dom.insert(1, ('order_id.team_id', 'in', tids))
if team_id is not None:
dom = [d for d in dom if not (isinstance(d, tuple) and d[0] == 'order_id.team_id')]
dom.append(('order_id.team_id', '=', team_id))
if date_from:
dom.append(('order_id.date_order', '>=', f'{date_from} 00:00:00'))
if date_to:
dom.append(('order_id.date_order', '<=', f'{date_to} 23:59:59'))
ex = excluded_partner_ids()
if ex:
dom.append(('order_partner_id', 'not in', list(ex)))
if partner_ids is not None:
dom.append(('order_partner_id', 'in', list(partner_ids)))
if doc_mode() == 'invoice':
dom.append(('order_id.invoice_status', '=', 'invoiced'))
if extra:
dom.extend(extra)
return dom
def read_group(model, domain, fields, groupby, **kw):
return get_odoo().read_group(model, domain=domain, fields=fields, groupby=groupby, **kw)
def search_read(model, domain=None, fields=None, **kw):
return get_odoo().search_read(model, domain=domain or [], fields=fields or [], **kw)
def sum_field(model, domain, field):
"""One-line aggregate sum of `field` over a domain (uses read_group, no row fetch).
Robust to the empty match set: a groupby=[] read_group over zero rows makes Odoo return a
None aggregate, which the server's XML-RPC layer cannot marshal ("cannot marshal None").
An empty set simply means the sum is 0.0, so we treat that specific Fault as zero."""
try:
g = get_odoo().read_group(model, domain=domain, fields=[f'{field}:sum'], groupby=[], lazy=False)
except Exception as e: # noqa: BLE001 — narrow check below
if 'cannot marshal None' in str(e):
return 0.0
raise
return (g[0].get(field) or 0.0) if g else 0.0
def distinct_count(model, domain, field):
"""Distinct count of a (m2o) field over a domain via grouped read_group."""
g = get_odoo().read_group(model, domain=domain, fields=[field], groupby=[field], lazy=False)
return sum(1 for r in g if r.get(field))