File size: 20,850 Bytes
c14ceee
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
c3e4cb4
 
 
 
 
 
051f280
 
 
 
 
 
 
 
 
 
 
 
 
c14ceee
051f280
 
 
 
 
 
 
c14ceee
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
051f280
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
609fb78
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
051f280
 
609fb78
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
c14ceee
609fb78
c14ceee
 
 
 
37d92f0
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
051f280
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
c14ceee
 
 
 
 
 
 
 
 
051f280
 
 
 
 
 
 
 
 
 
 
 
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
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
"""core/keychain.py β€” per-tenant CREDENTIAL STORE (wave 18, contract C7 / ruling R3).

An admin stores a data source's secrets (Royal's Odoo API key today; any tenant's sources
tomorrow) encrypted at rest in the TENANT's store bucket. The runtime posture is R3's
"live with env fallback": the connector layer asks `odoo_creds(rt)` FIRST and falls back to
the process environment (tenant #0's `.env` / Space secrets), so Royal keeps working with an
empty keychain and a new tenant needs no code to bring keys. The FULL cutover (env retired)
is staged for a later wave β€” R3, verbatim.

CRYPTO: Fernet (AES128-CBC + HMAC, the `cryptography` package) under ONE platform key,
`AIOS_KEYCHAIN_KEY` (a `Fernet.generate_key()` value; a Space secret, never in the store).
No key β‡’ the keychain is LOCKED: metadata still lists (labels are not secrets), but nothing
decrypts and writes are refused β€” fail closed, never a plaintext fallback.

SHAPE (bucket `keychain`, per-tenant via the runtime handle):

    {entries: {id: {id, label, type, enc, preview, created, createdBy}}}

`enc` = Fernet token over the JSON fields dict. `preview` = a masked hint ("Β·Β·Β·Β·1234")
computed ONCE at write from the entry's most secret-looking field β€” the ONLY part of a
secret that ever leaves this module through a route. `read_fields` exists for the CONNECTOR
LAYER; no route returns its output.
"""
import datetime as _dt
import json
import os
import secrets as _secrets

KEY = 'keychain'
#: ⭐ WAVE 23 (C11/R8) β€” `stripe` and `shopify` joined. They are TOKEN connectors, not OAuth
#: ones, and that is the whole reason they could ship in this wave: both issue a long-lived
#: restricted key / custom-app token an admin pastes in, so they need no OAuth client, no
#: consent screen and no app review β€” the three things that make every other provider on the
#: standing queue a multi-week errand. A type here is what lets the connectors directory show
#: them as CONNECTABLE rather than as another faded to-do.
#: ⭐⭐ WAVE 31 (R2/R5) β€” `meta_ads` JOINS, AND IT IS THE FIRST END-TO-END KEYCHAINβ†’HTTP CONNECTOR.
#: Every type above it is either read by nothing (`generic`, `stripe`, `shopify` store a key that
#: no data path consumes yet) or read by ONE resolver that talks XML-RPC (`odoo`). So `meta_ads`
#: is the first entry whose stored secret actually reaches an HTTP API in production, through
#: `meta_creds` below β€” which is why W31-T48's `how:` says *build the seam, do not assume it*.
#: β›” A TYPE MISSING FROM THIS SET 400s IN `add_entry`, which is the entire failure: the Connectors
#: directory would offer Meta Ads, an admin would paste a token, and the write would be refused
#: with "type must be one of odoo, generic, stripe, shopify" β€” a message about our data structure,
#: on a screen about their advertising account.
#: ⚠ R8 keeps the two doors separate on purpose: this is the TENANT's own token (GTM Lab pastes
#: theirs), while R5's "Connect Facebook" is the delegated-OAuth door Nango owns. Neither replaces
#: the other, and the schema work is blocked on neither.
ENTRY_TYPES = {'odoo', 'generic', 'stripe', 'shopify', 'meta_ads'}
#: which stored field feeds the masked preview, per type (first present wins)
#: ⚠ `access_token` JOINED IN WAVE 31, AND ITS ABSENCE WAS A REAL DEFECT FOUND BY BUILDING THE
#: SEAM RATHER THAN BY READING IT. Meta's credential is conventionally an `access_token`, and with
#: only `token` on this list `add_entry` stored the secret correctly and computed an EMPTY preview
#: β€” so the keychain table would show a Meta row with no hint at all, and an admin holding two
#: tokens could not tell which one is stored or whether the paste worked. Nothing raises; the
#: entry is fine; the only symptom is a blank cell on the screen the feature exists for.
_PREVIEW_FIELDS = ('api_key', 'access_token', 'value', 'password', 'token')
MAX_ENTRIES = 40
MAX_FIELD_LEN = 500


class KeychainLocked(Exception):
    """Raised when AIOS_KEYCHAIN_KEY is absent/wrong β€” callers answer 503, never {}."""


def _fernet():
    key = (os.environ.get('AIOS_KEYCHAIN_KEY') or '').strip()
    if not key:
        raise KeychainLocked('AIOS_KEYCHAIN_KEY is not configured')
    try:
        from cryptography.fernet import Fernet
        return Fernet(key.encode('ascii'))
    except Exception as e:
        raise KeychainLocked(f'keychain key unusable: {type(e).__name__}')


def unlocked():
    try:
        _fernet()
        return True
    except KeychainLocked:
        return False


def _bucket(rt):
    return rt.get(KEY) or {}


def list_entries(rt):
    """Metadata ONLY β€” never a decrypted field, never `enc` itself."""
    out = []
    for eid, e in sorted((_bucket(rt).get('entries') or {}).items()):
        if not isinstance(e, dict):
            continue
        out.append({'id': eid, 'label': e.get('label') or eid,
                    'type': e.get('type') or 'generic',
                    'preview': e.get('preview') or '',
                    'created': e.get('created') or '',
                    'createdBy': e.get('createdBy') or ''})
    return out


def add_entry(rt, label, etype, fields, username):
    """Encrypt + store. Returns the metadata row, or raises ValueError on shape problems and
    KeychainLocked when there is no key (a secret must never be stored in the clear)."""
    label = ' '.join(str(label or '').split())[:80]
    etype = str(etype or '').strip().lower()
    if not label:
        raise ValueError('give the key a label')
    if etype not in ENTRY_TYPES:
        raise ValueError(f"type must be one of {', '.join(sorted(ENTRY_TYPES))}")
    if not isinstance(fields, dict) or not fields:
        raise ValueError('fields must be a non-empty object')
    clean = {}
    for k, v in list(fields.items())[:12]:
        k = str(k).strip()[:40]
        if k and isinstance(v, (str, int, float, bool)):
            clean[k] = str(v)[:MAX_FIELD_LEN]
    if not clean:
        raise ValueError('no usable fields')
    token = _fernet().encrypt(json.dumps(clean, ensure_ascii=False).encode('utf-8'))
    secretish = next((clean[k] for k in _PREVIEW_FIELDS if clean.get(k)), '')
    preview = ('Β·Β·Β·Β·' + secretish[-4:]) if len(secretish) >= 4 else ('Β·Β·Β·Β·' if secretish else '')
    eid = f'k_{_secrets.token_hex(6)}'
    row = {'id': eid, 'label': label, 'type': etype, 'enc': token.decode('ascii'),
           'preview': preview,
           'created': _dt.datetime.now().strftime('%Y-%m-%dT%H:%M:%S'),
           'createdBy': str(username or '')}

    def _up(cur):
        entries = cur.setdefault('entries', {})
        if len(entries) >= MAX_ENTRIES:
            raise ValueError(f'this keychain is at its {MAX_ENTRIES}-entry cap')
        entries[eid] = row
        return cur

    rt.update(KEY, _up, flush='sync')
    return {'id': eid, 'label': label, 'type': etype, 'preview': preview,
            'created': row['created'], 'createdBy': row['createdBy']}


def delete_entry(rt, entry_id):
    def _drop(cur):
        (cur.get('entries') or {}).pop(str(entry_id), None)
        return cur
    rt.update(KEY, _drop, flush='sync')
    return True


def read_fields(rt, entry_id):
    """Decrypt ONE entry's fields β€” for the CONNECTOR LAYER only; no route returns this.
    None when the entry does not exist; KeychainLocked when the key is absent/wrong (a wrong
    key must be loud β€” silently empty creds would read as 'not configured')."""
    e = (_bucket(rt).get('entries') or {}).get(str(entry_id))
    if not isinstance(e, dict) or not e.get('enc'):
        return None
    try:
        raw = _fernet().decrypt(str(e['enc']).encode('ascii'))
    except KeychainLocked:
        raise
    except Exception:
        raise KeychainLocked('this entry does not decrypt under the configured key')
    try:
        data = json.loads(raw.decode('utf-8'))
        return data if isinstance(data, dict) else None
    except Exception:
        return None


def odoo_creds(rt):
    """R3's resolver seam: the FIRST odoo-type entry's fields, else None (the caller falls
    back to the environment). Deterministic order = insertion-id sort, so 'first' is stable."""
    return _first_creds(rt, 'odoo')


def meta_creds(rt):
    """W31 (R2): the FIRST `meta_ads` entry's fields, else None β€” the Meta connector's resolver.

    β›” IT IS DELIBERATELY THE SAME SHAPE AS `odoo_creds` AND SHARES ITS IMPLEMENTATION, because
    "which stored credential serves this connector?" is ONE question and this module was one
    copy-paste away from having two answers to it ([[one-question-two-normalizers]]). The
    difference between the two connectors belongs in the CONNECTOR, not in the lookup.

    ⚠ AND IT FAILS CLOSED WITH NO ENVIRONMENT FALLBACK, which is where it deliberately DIFFERS
    from `odoo_creds`'s caller. `harness/runtime.py::odoo_source` falls back to the process
    environment for tenant #0 β€” a documented exception for the tenant whose `.env` this is. There
    is no equivalent for Meta: `META_ADS_ACCESS_TOKEN` is the OWNER's measurement token (R8), and
    handing it to a tenant whose admin has not stored one is precisely the leak the R3 rule
    exists to prevent. A tenant with no entry gets None, and the connector says so.
    """
    return _first_creds(rt, 'meta_ads')


#: ⭐⭐ WAVE 32 (R4 / contract C1) β€” THE SCOPE SIDE-BUCKET, READ HERE, WRITTEN NOWHERE IN THIS FILE.
#:
#: `{entry_id: {"scope": "personal", "owner": "<username>"}}`. A `business` entry writes NO row, so
#: an absent row and a `business` row mean the same thing β€” which is why this is a READ DEFAULT and
#: never a migration ([[a-migration-that-runs-on-the-next-write]]).
#:
#: β›” DECLARED TWICE ON PURPOSE, AND GATED. The vocabulary's home is `api/routes_keychain.py`
#: (SESSION B's file, C1 is explicit that it is declared there). `core/` may not import from `api/`
#: β€” the layer contract points one way β€” so this reads the same store key by name. Two spellings of
#: one constant is [[a-constant-two-features-share]], so `verify_api` asserts the two agree, with an
#: NC. It is the same posture `store.backend()` takes toward `store_backend.name()`.
SCOPE_KEY = 'keychain_scopes'

#: The credential types a WHOLE WORKSPACE reads through, which therefore cannot be personal.
#: β›” THE ELEVATION THIS PREVENTS, in B's words: `odoo_creds`/`meta_creds` resolve to *the first
#: entry of that type for the TENANT* β€” that is what spawns every `ut_odoo_*` / `ut_meta_*` grid and
#: what every measure column is answered from. So a member storing a personal Odoo key would not get
#: "their own Odoo": they would silently become the credential the entire workspace's databases are
#: built from. **A credential elevation wearing a scope picker.**
TENANT_WIDE_TYPES = ('odoo', 'meta_ads')


def _personal_ids(rt):
    """Entry ids marked `personal`. Absent bucket β‡’ empty β‡’ today's behaviour, unchanged."""
    try:
        rows = rt.get(SCOPE_KEY) or {}
    except Exception:                                                     # noqa: BLE001
        return frozenset()
    if not isinstance(rows, dict):
        return frozenset()
    return frozenset(str(k) for k, v in rows.items()
                     if isinstance(v, dict) and str(v.get('scope') or '') == 'personal')


def _first_creds(rt, etype):
    """The first entry of `etype`, decrypted β€” or None. Deterministic order = insertion-id sort,
    so "first" is stable across reads rather than dict-order luck.

    ⭐⭐ WAVE 32 (R4/C1) β€” A PERSONAL ENTRY OF A TENANT-WIDE TYPE IS SKIPPED, NOT SERVED.
    `routes_keychain` REFUSES to create one, and that closes the door from the only side B owns;
    this is the other side, and it is the one that matters for anything already stored β€” an entry
    written before this wave, by an older client, or by any path that does not go through that
    route. B booked it for A explicitly ("the resolver-side guard is booked for A").

    ⚠ SKIP, NOT RAISE. The caller's contract is "the tenant's credential, or None", and None already
    means "not configured" everywhere it is consumed: `odoo_source` falls back to the environment
    for tenant #0, and Meta fails closed with no fallback. Raising would turn a mis-scoped entry
    into an outage for a workspace that has a perfectly good business entry two rows down.
    ⚠ AND THE SKIP IS SCOPED TO `TENANT_WIDE_TYPES`. A personal entry of any OTHER type is a normal
    personal credential and this resolver is not how it is reached.
    """
    skip = _personal_ids(rt) if etype in TENANT_WIDE_TYPES else frozenset()
    for row in list_entries(rt):
        if row['type'] == etype and str(row['id']) not in skip:
            return read_fields(rt, row['id'])
    return None


# ---------------------------------------------------------------------------------------------
# PER-USER SECRET SLOTS (wave 22, contract C5 β€” the R7 seam adopted from QM)
# ---------------------------------------------------------------------------------------------
# A connector credential that belongs to a PERSON, not to the tenant: *my* Gmail refresh token,
# never "the workspace's". Stored in a SEPARATE sub-bucket (`user_secrets`) beside `entries`,
# deliberately: `list_entries` walks `entries` only, so a user's OAuth tokens never show up in
# the admin keychain table β€” they are identity, not shared infrastructure β€” and the addition is
# migration-free (existing entries and their readers are untouched). Same Fernet, same
# fail-closed law: no AIOS_KEYCHAIN_KEY β‡’ writes refuse, reads raise, nothing plaintext.

#: Slot names are SHAPE-validated rather than enumerated (wave 22, C5 amendment A2): the OAuth
#: provider registry lives API-side and grows one entry per provider β€” a platform-side copy of
#: its keys would be the client-union defect one layer down. The shape is the law instead.
import re as _re_slots

USER_SECRET_SLOT = _re_slots.compile(r'^oauth_[a-z0-9_]{1,24}$')


def put_user_secret(rt, username, provider, fields):
    """Encrypt + store ONE user's credential slot for `provider` (upsert β€” a reconnect
    replaces). Raises ValueError on shape problems, KeychainLocked when there is no key."""
    username = str(username or '').strip()
    provider = str(provider or '').strip().lower()
    if not username:
        raise ValueError('no username for the secret slot')
    if not USER_SECRET_SLOT.fullmatch(provider):
        raise ValueError('a secret slot is named oauth_<provider>')
    if not isinstance(fields, dict) or not fields:
        raise ValueError('fields must be a non-empty object')
    clean = {}
    for k, v in list(fields.items())[:16]:
        k = str(k).strip()[:40]
        if k and isinstance(v, (str, int, float, bool)):
            clean[k] = str(v)[:2000]         # a Google refresh token is ~100 chars; JWTs more
    if not clean:
        raise ValueError('no usable fields')
    token = _fernet().encrypt(json.dumps(clean, ensure_ascii=False).encode('utf-8'))
    row = {'enc': token.decode('ascii'),
           'updated': _dt.datetime.now().strftime('%Y-%m-%dT%H:%M:%S')}

    def _up(cur):
        cur.setdefault('user_secrets', {}).setdefault(username, {})[provider] = row
        return cur

    rt.update(KEY, _up, flush='sync')
    return True


def read_user_secret(rt, username, provider):
    """Decrypt ONE user's slot β€” None when absent; KeychainLocked stays LOUD (a wrong key must
    never read as 'not connected')."""
    slot = ((_bucket(rt).get('user_secrets') or {}).get(str(username)) or {}).get(
        str(provider or '').lower())
    if not isinstance(slot, dict) or not slot.get('enc'):
        return None
    try:
        raw = _fernet().decrypt(str(slot['enc']).encode('ascii'))
    except KeychainLocked:
        raise
    except Exception:
        raise KeychainLocked('this slot does not decrypt under the configured key')
    try:
        data = json.loads(raw.decode('utf-8'))
        return data if isinstance(data, dict) else None
    except Exception:
        return None


def drop_user_secret(rt, username, provider):
    def _drop(cur):
        slots = (cur.get('user_secrets') or {}).get(str(username))
        if isinstance(slots, dict):
            slots.pop(str(provider or '').lower(), None)
        return cur
    rt.update(KEY, _drop, flush='sync')
    return True


# ═════════════════════════════════════════════════════════════════════════════════════════════
# ⭐⭐ W31 β€” THE LIVENESS-PROBE SLOT, AND WHY IT IS A SLOT AND NOT AN IMPORT.
# ═════════════════════════════════════════════════════════════════════════════════════════════
#
# `test_entry` below wants to ASK each connector's API whether a stored credential works. The
# knowledge of HOW to ask lives in the connector modules β€” and those live in the API layer, which
# imports `core`, never the other way round (ARCHITECTURE.md's layering rule, restated in
# CLAUDE.md: *"core never imports up"*). A first draft of the Meta arm did the tempting thing β€”
# `sys.path.insert` into `aios-web/api` and `import connectors_meta` from inside `core` β€” and that
# is an architecture inversion with no gate on it, i.e. the kind that survives.
#
# β›” THE FIX IS THIS REPO'S OWN IDIOM, NOT A NEW ONE. `harness/datastore.py` has exactly this
# problem with the connector-pause flag and solves it with `set_paused_probe(fn)`: the lower layer
# exposes a SLOT and the upper layer fills it at import, with the comment *"an import back would
# be a cycle"*. Same shape, same reason. `verify_meta` asserts `core/` imports nothing from the
# API layer, so the inversion cannot come back quietly.
#
#: entry type -> `fn(fields: dict, timeout: int) -> {ok, message}`. Filled by the connector layer.
_PROBES = {}


def register_prober(etype, fn):
    """Let the connector layer answer `test_entry` for ONE entry type.

    ⚠ Additive and idempotent by design β€” a re-import must not double-register or raise, because
    module import order is not something a caller controls.
    """
    etype = str(etype or '').strip().lower()
    if etype and callable(fn):
        _PROBES[etype] = fn
    return etype in _PROBES


def test_entry(rt, entry_id, timeout=8):
    """A cheap, READ-ONLY liveness probe. odoo: xmlrpc `common.version()` (no auth, no data);
    generic: shape-only. Always answers {ok, message} β€” an exception here is a result, not a
    crash."""
    try:
        e = (_bucket(rt).get('entries') or {}).get(str(entry_id))
        if not isinstance(e, dict):
            return {'ok': False, 'message': 'no such entry'}
        fields = read_fields(rt, entry_id) or {}
        # ⭐ W31 β€” A REGISTERED PROBER ANSWERS FOR ITS OWN TYPE, if the connector layer filled the
        # slot. A shape check says "you pasted something", which is the least useful thing an
        # admin can be told about a credential: a revoked or wrong-scoped token is the same SHAPE
        # as a working one, so a connector that can ASK its API should.
        prober = _PROBES.get(str(e.get('type') or ''))
        if prober is not None:
            try:
                got = prober(dict(fields), timeout)
                return got if isinstance(got, dict) else {'ok': False,
                                                          'message': 'the prober said nothing'}
            except Exception as ex:                              # noqa: BLE001
                return {'ok': False, 'message': f'{type(ex).__name__}: {str(ex)[:160]}'}
        if e.get('type') != 'odoo':
            return {'ok': bool(fields), 'message': f'{len(fields)} field(s) stored'}
        url = str(fields.get('url') or '').rstrip('/')
        if not url.startswith(('http://', 'https://')):
            return {'ok': False, 'message': 'url must start with http(s)://'}
        import socket
        import xmlrpc.client
        old = socket.getdefaulttimeout()
        socket.setdefaulttimeout(timeout)
        try:
            common = xmlrpc.client.ServerProxy(f'{url}/xmlrpc/2/common', allow_none=True)
            ver = common.version()
        finally:
            socket.setdefaulttimeout(old)
        sv = (ver or {}).get('server_version') if isinstance(ver, dict) else None
        return {'ok': True, 'message': f'Odoo answered (server {sv or "unknown"})'}
    except KeychainLocked as e:
        return {'ok': False, 'message': str(e)}
    except Exception as e:
        return {'ok': False, 'message': f'{type(e).__name__}: {str(e)[:120]}'}