File size: 13,443 Bytes
c14ceee
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
0adc1e7
c14ceee
 
 
 
 
 
 
 
 
 
0adc1e7
 
 
 
c14ceee
 
 
 
 
 
 
0adc1e7
 
c14ceee
 
 
 
 
 
 
 
 
 
 
 
 
 
0adc1e7
c14ceee
 
 
0adc1e7
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
c14ceee
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
0adc1e7
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
c14ceee
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
0adc1e7
 
c14ceee
 
 
 
0adc1e7
c14ceee
 
 
 
 
0adc1e7
c14ceee
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
0adc1e7
 
 
 
c14ceee
0adc1e7
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
c14ceee
 
 
 
 
 
 
0adc1e7
 
 
 
c14ceee
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
"""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))