From 35297ff79f22b5d8dfb482e5781576025dd56e23 Mon Sep 17 00:00:00 2001 From: dada-yan Date: Mon, 21 Sep 2026 00:38:08 +1000 Subject: [PATCH 1/5] Model action dispatch identity and uncertain outcomes before enabling a sender Signed-off-by: dada-yan --- scripts/prototypes/actions/protocol.py | 153 +++++++++++++ scripts/prototypes/actions/test_protocol.py | 228 ++++++++++++++++++++ 2 files changed, 381 insertions(+) create mode 100644 scripts/prototypes/actions/protocol.py create mode 100644 scripts/prototypes/actions/test_protocol.py diff --git a/scripts/prototypes/actions/protocol.py b/scripts/prototypes/actions/protocol.py new file mode 100644 index 000000000..df6e3da5f --- /dev/null +++ b/scripts/prototypes/actions/protocol.py @@ -0,0 +1,153 @@ +"""Isolated PostgreSQL protocol experiment, NOT Utopia's action implementation. + +No production tables, routes, credentials or generic worker integration. Targets +are a test server on loopback only. The model establishes dispatch identities and +fault boundaries; it does not establish Utopia RBAC, templates or DNS pinning. +""" +import hashlib +import http.client +import json +import os +import uuid +from contextlib import contextmanager + +import psycopg +from psycopg import sql +from psycopg.rows import dict_row + + +class Conflict(Exception): + pass + + +class Model: + def __init__(self, dsn, port): + self.dsn, self.port = dsn, port + self.schema = 'action_model_' + uuid.uuid4().hex + with psycopg.connect(dsn, autocommit=True) as c: + c.execute(sql.SQL('CREATE SCHEMA {}').format(sql.Identifier(self.schema))) + with self.connect() as c: + c.execute('''CREATE TABLE definition ( + id integer PRIMARY KEY, revision integer NOT NULL, + enabled boolean NOT NULL, granted boolean NOT NULL, + editor boolean NOT NULL); + INSERT INTO definition VALUES(1,1,true,true,true); + CREATE TABLE runs ( + id uuid PRIMARY KEY, actor text NOT NULL, scope uuid, + request_id uuid NOT NULL, fingerprint text NOT NULL, + revision integer NOT NULL, state text NOT NULL, + token uuid, status integer, capture text, excerpt text, + UNIQUE NULLS NOT DISTINCT(actor, scope, request_id));''') + + @contextmanager + def connect(self): + with psycopg.connect(self.dsn, row_factory=dict_row) as c: + c.execute(sql.SQL('SET search_path TO {}').format(sql.Identifier(self.schema))) + yield c + + def close(self): + with psycopg.connect(self.dsn, autocommit=True) as c: + c.execute(sql.SQL('DROP SCHEMA {} CASCADE').format(sql.Identifier(self.schema))) + + def definition(self, **changes): + with self.connect() as c: + for key, value in changes.items(): + assert key in ('revision', 'enabled', 'granted', 'editor') + c.execute(sql.SQL('UPDATE definition SET {}=%s').format(sql.Identifier(key)), (value,)) + + @staticmethod + def authorize(row, revision): + if not row['enabled'] or not row['granted'] or not row['editor']: + raise PermissionError('denied at dispatch boundary') + if row['revision'] != revision: + raise Conflict('preview_stale') + + def preview(self, args): + with self.connect() as c: + d = c.execute('SELECT * FROM definition WHERE id=1').fetchone() + self.authorize(d, d['revision']) + return {'revision': d['revision'], 'args': args} + + def prepare(self, request_id, args, revision=1, scope=None, actor='editor'): + fingerprint = hashlib.sha256(json.dumps([revision, args], sort_keys=True, + separators=(',', ':'), allow_nan=False).encode()).hexdigest() + with self.connect() as c: + # Replays are authorized as well. This model has one definition, not + # production role tables; it tests the cutoff, not the RBAC resolver. + d = c.execute('SELECT * FROM definition WHERE id=1 FOR SHARE').fetchone() + self.authorize(d, revision) + new = c.execute('''INSERT INTO runs(id,actor,scope,request_id,fingerprint,revision,state) + VALUES(%s,%s,%s,%s,%s,%s,'prepared') ON CONFLICT DO NOTHING RETURNING *''', + (uuid.uuid4(), actor, scope, request_id, fingerprint, revision)).fetchone() + if new: + return new, True + old = c.execute('SELECT * FROM runs WHERE actor=%s AND scope IS NOT DISTINCT FROM %s AND request_id=%s', + (actor, scope, request_id)).fetchone() + if old['fingerprint'] != fingerprint: + raise Conflict('execution_request_reused') + return old, False + + def gate(self, run_id, revision, commit_uncertain=False): + token = uuid.uuid4() + with self.connect() as c: + d = c.execute('SELECT * FROM definition WHERE id=1 FOR SHARE').fetchone() + self.authorize(d, revision) + got = c.execute("UPDATE runs SET state='dispatching',token=%s WHERE id=%s AND state='prepared' RETURNING id", + (token, run_id)).fetchone() + # Commit happened, but caller cannot know: it must NOT send. Deliberate + # fault injection after real commit models a lost commit acknowledgement. + if commit_uncertain: + raise ConnectionError('injected lost gate commit acknowledgement') + return token if got else None + + def read(self, run_id): + with self.connect() as c: + return c.execute('SELECT * FROM runs WHERE id=%s', (run_id,)).fetchone() + + def expire(self, run_id): + with self.connect() as c: + c.execute("UPDATE runs SET state='not_sent' WHERE id=%s AND state='prepared'", (run_id,)) + + def recover(self): + with self.connect() as c: + c.execute("UPDATE runs SET state='outcome_unknown' WHERE state='dispatching'") + + def observe(self, run_id, token, state, status=None, capture='absent', excerpt=''): + with self.connect() as c: + c.execute('''UPDATE runs SET state=%s,status=%s,capture=%s,excerpt=%s + WHERE id=%s AND token=%s AND state IN ('dispatching','outcome_unknown')''', + (state, status, capture, excerpt, run_id, token)) + + def send(self, run_id, token, path, fail_write=False): + # Fixed loopback destination: this is intentionally NOT a reusable sender. + # http.client neither follows redirects nor retries requests. No proxies. + conn = http.client.HTTPConnection('127.0.0.1', self.port, timeout=0.3) + status, capture, excerpt = None, 'absent', '' + try: + conn.request('POST', path, body=b'{}', headers={'Content-Type': 'application/json'}) + response = conn.getresponse() + status = response.status + try: + body = response.read(1025) + capture = 'truncated' if len(body) > 1024 else 'complete' + excerpt = body[:256].decode('utf-8', errors='replace') + except (OSError, http.client.HTTPException): + capture = 'read_error' + except (OSError, http.client.HTTPException): + pass + finally: + conn.close() + if fail_write: + raise ConnectionError('injected observation persistence failure') + self.observe(run_id, token, 'response_received' if status is not None else 'outcome_unknown', + status, capture, excerpt) + + def execute(self, request_id, args, revision=1, scope=None, path='/ok'): + row, created = self.prepare(request_id, args, revision, scope) + # Only the original creator proceeds. Returning prepared after a restart + # does not grant replay callers permission to take over and send. + if created: + token = self.gate(row['id'], revision) + if token: + self.send(row['id'], token, path) + return self.read(row['id']) diff --git a/scripts/prototypes/actions/test_protocol.py b/scripts/prototypes/actions/test_protocol.py new file mode 100644 index 000000000..e35ea0df1 --- /dev/null +++ b/scripts/prototypes/actions/test_protocol.py @@ -0,0 +1,228 @@ +import concurrent.futures +import http.server +import os +import socket +import threading +import unittest +import uuid +from protocol import Model, Conflict + + +class Remote(http.server.BaseHTTPRequestHandler): + def log_message(self, *_): + pass + + def do_POST(self): + self.rfile.read(int(self.headers.get('Content-Length', 0))) + with self.server.guard: + self.server.requests.append(self.path) + self.server.effects += 1 + if self.path == '/drop': + # A real remote effect is committed BEFORE losing the response. + self.server.effect_committed.set() + self.server.release.wait(5) + self.connection.shutdown(socket.SHUT_RDWR) + self.connection.close() + return + if self.path == '/slow-body': + self.send_response(202) + self.send_header('Content-Length', '100') + self.end_headers() + self.wfile.flush() + self.server.release.wait(2) + return + status = int(self.path[1:]) if self.path[1:].isdigit() else 200 + body = ('汉字' * 1000).encode() if self.path == '/large' else b'observed response' + self.send_response(status) + if status == 302: + self.send_header('Location', '/second-target') + self.send_header('Content-Length', str(len(body))) + self.end_headers() + self.wfile.write(body) + + +class ProtocolTests(unittest.TestCase): + def setUp(self): + self.remote = http.server.ThreadingHTTPServer(('127.0.0.1', 0), Remote) + self.remote.requests, self.remote.effects = [], 0 + self.remote.guard = threading.Lock() + self.remote.effect_committed = threading.Event() + self.remote.release = threading.Event() + self.thread = threading.Thread(target=self.remote.serve_forever, daemon=True) + self.thread.start() + self.model = Model(os.environ['UTOPIA_DATABASE_URL'], self.remote.server_port) + self.key = uuid.uuid4() + + def tearDown(self): + self.remote.release.set() + self.remote.shutdown() + self.remote.server_close() + self.thread.join() + self.model.close() + + def prepared(self): + return self.model.prepare(self.key, {'n': 1})[0] + + def test_preview_has_neither_run_nor_network_effect(self): + self.assertEqual(self.model.preview({'n': 1})['revision'], 1) + with self.model.connect() as c: + self.assertEqual(c.execute('SELECT count(*) n FROM runs').fetchone()['n'], 0) + self.assertEqual(self.remote.effects, 0) + + def test_stale_preview_and_revoked_authority_never_dispatch(self): + self.model.definition(revision=2) + with self.assertRaises(Conflict): + self.model.execute(self.key, {}, revision=1) + self.model.definition(revision=1) + row = self.prepared() + for field in ['enabled', 'granted', 'editor']: + self.model.definition(**{field: False}) + with self.assertRaises(PermissionError): + self.model.gate(row['id'], 1) + with self.assertRaises(PermissionError): + self.model.execute(self.key, {'n': 1}) + self.model.definition(**{field: True}) + self.assertEqual(self.remote.effects, 0) + + def test_concurrent_duplicate_registry_scope_has_one_run_and_one_dispatch(self): + barrier = threading.Barrier(8) + def call(_): + barrier.wait(timeout=5) + return self.model.execute(self.key, {'n': 1})['id'] + with concurrent.futures.ThreadPoolExecutor(8) as executor: + ids = list(executor.map(call, range(8))) + self.assertEqual(len(set(ids)), 1) + self.assertEqual(self.remote.effects, 1) + + def test_kb_scope_retries_and_new_explicit_operation(self): + scope = uuid.uuid4() + a = self.model.execute(self.key, {'n': 1}, scope=scope) + b = self.model.execute(self.key, {'n': 1}, scope=scope) + self.assertEqual(a['id'], b['id']) + self.model.execute(uuid.uuid4(), {'n': 1}, scope=scope) + self.assertEqual(self.remote.effects, 2) + + def test_same_identity_cannot_change_inputs(self): + self.model.execute(self.key, {'n': 1}) + with self.assertRaises(Conflict): + self.model.execute(self.key, {'n': 2}) + self.assertEqual(self.remote.effects, 1) + + def test_prepared_insert_failure_cannot_send(self): + with self.model.connect() as c: + c.execute("ALTER TABLE runs ADD CONSTRAINT injected_failure CHECK (state <> 'prepared')") + with self.assertRaises(Exception): + self.model.execute(self.key, {}) + self.assertEqual(self.remote.effects, 0) + + def test_prepared_replay_is_not_a_recovery_sender(self): + row = self.prepared() + self.assertEqual(self.model.execute(self.key, {'n': 1})['state'], 'prepared') + self.model.expire(row['id']) + self.assertIsNone(self.model.gate(row['id'], 1)) + self.assertEqual(self.remote.effects, 0) + + def test_expiration_and_dispatch_gate_are_mutually_exclusive(self): + row = self.prepared() + barrier = threading.Barrier(2) + def expire(): + barrier.wait(timeout=5) + self.model.expire(row['id']) + def gate(): + barrier.wait(timeout=5) + return self.model.gate(row['id'], 1) + with concurrent.futures.ThreadPoolExecutor(2) as ex: + a, b = ex.submit(expire), ex.submit(gate) + a.result() + token = b.result() + self.assertEqual(self.model.read(row['id'])['state'], 'dispatching' if token else 'not_sent') + self.assertIsNone(self.model.gate(row['id'], 1)) + + def test_lost_gate_commit_acknowledgement_does_not_authorize_send(self): + row = self.prepared() + with self.assertRaises(ConnectionError): + self.model.gate(row['id'], 1, commit_uncertain=True) + self.assertEqual(self.model.read(row['id'])['state'], 'dispatching') + self.model.recover() + self.assertEqual(self.model.execute(self.key, {'n': 1})['state'], 'outcome_unknown') + self.assertEqual(self.remote.effects, 0) + + def test_gate_then_crash_before_send_is_conservatively_unknown(self): + row = self.prepared() + self.model.gate(row['id'], 1) + self.model.recover() + self.assertEqual(self.model.execute(self.key, {'n': 1})['state'], 'outcome_unknown') + self.assertEqual(self.remote.effects, 0) + + def test_remote_effect_then_dropped_response_is_not_retried(self): + with concurrent.futures.ThreadPoolExecutor(1) as ex: + future = ex.submit(self.model.execute, self.key, {'n': 1}, path='/drop') + try: + self.assertTrue(self.remote.effect_committed.wait(3)) + self.assertEqual(self.remote.effects, 1) + finally: + self.remote.release.set() + row = future.result(timeout=3) + self.assertEqual(row['state'], 'outcome_unknown') + self.model.recover() + self.model.execute(self.key, {'n': 1}) + self.assertEqual(self.remote.effects, 1) + + def test_http_status_is_observation_not_business_success(self): + for status in [200, 202, 400, 500]: + row = self.model.execute(uuid.uuid4(), {}, path=f'/{status}') + self.assertEqual(row['state'], 'response_received') + self.assertEqual(row['status'], status) + self.assertEqual(row['capture'], 'complete') + self.assertEqual(self.remote.effects, 4) + + def test_body_timeout_preserves_observed_status(self): + row = self.model.execute(self.key, {}, path='/slow-body') + self.assertEqual((row['state'], row['status'], row['capture']), ('response_received', 202, 'read_error')) + self.assertEqual(self.remote.effects, 1) + + def test_body_cap_is_reported_and_excerpt_is_unicode(self): + row = self.model.execute(self.key, {}, path='/large') + self.assertEqual(row['capture'], 'truncated') + self.assertIsInstance(row['excerpt'], str) + self.assertLessEqual(len(row['excerpt']), 256) + self.assertEqual(self.remote.effects, 1) + + def test_observation_write_failure_recovery_never_resends(self): + row = self.prepared() + token = self.model.gate(row['id'], 1) + with self.assertRaises(ConnectionError): + self.model.send(row['id'], token, '/ok', fail_write=True) + self.model.recover() + self.assertEqual(self.model.execute(self.key, {'n': 1})['state'], 'outcome_unknown') + self.assertEqual(self.remote.effects, 1) + + def test_late_observation_uses_original_token_without_dispatch(self): + row = self.prepared() + token = self.model.gate(row['id'], 1) + self.model.recover() + self.model.observe(row['id'], uuid.uuid4(), 'response_received', 200) + self.assertEqual(self.model.read(row['id'])['state'], 'outcome_unknown') + self.model.observe(row['id'], token, 'response_received', 202) + self.assertEqual(self.model.execute(self.key, {'n': 1})['status'], 202) + self.assertEqual(self.remote.effects, 0) + + def test_redirect_is_an_observed_response_without_a_second_request(self): + row = self.model.execute(self.key, {}, path='/302') + self.assertEqual(row['status'], 302) + self.assertEqual(self.remote.requests, ['/302']) + + def test_proxy_environment_cannot_redirect_loopback_experiment(self): + old = os.environ.get('HTTP_PROXY') + os.environ['HTTP_PROXY'] = 'http://127.0.0.1:1' + try: + self.assertEqual(self.model.execute(self.key, {})['status'], 200) + finally: + if old is None: + os.environ.pop('HTTP_PROXY', None) + else: + os.environ['HTTP_PROXY'] = old + + +if __name__ == '__main__': + unittest.main(verbosity=2) From 297e66453857786efa1a6be7e32e33e90913b59d Mon Sep 17 00:00:00 2001 From: dada-yan Date: Mon, 21 Sep 2026 00:43:26 +1000 Subject: [PATCH 2/5] Kill action model processes before dispatch and after a remote effect Signed-off-by: dada-yan --- scripts/prototypes/actions/test_protocol.py | 41 +++++++++++++++++++++ 1 file changed, 41 insertions(+) diff --git a/scripts/prototypes/actions/test_protocol.py b/scripts/prototypes/actions/test_protocol.py index e35ea0df1..bf9dee5a1 100644 --- a/scripts/prototypes/actions/test_protocol.py +++ b/scripts/prototypes/actions/test_protocol.py @@ -223,6 +223,47 @@ def test_proxy_environment_cannot_redirect_loopback_experiment(self): else: os.environ['HTTP_PROXY'] = old + def test_real_process_exit_before_and_after_dispatch_boundary(self): + import subprocess + import sys + import selectors + for phase in ['prepared', 'dispatching', 'sending']: + key = uuid.uuid4() + child_code = """ +import os,sys,time +from protocol import Model +m=Model.__new__(Model) +m.dsn=os.environ['UTOPIA_DATABASE_URL'];m.schema=sys.argv[1];m.port=int(sys.argv[2]) +row,_=m.prepare(sys.argv[3], {}) +if sys.argv[4]!='prepared': token=m.gate(row['id'],1) +print('ready',flush=True) +if sys.argv[4]=='sending': m.send(row['id'],token,'/drop') +time.sleep(30) +""" + child = subprocess.Popen([sys.executable, '-c', child_code, self.model.schema, + str(self.remote.server_port), str(key), phase], + stdout=subprocess.PIPE, text=True) + try: + with selectors.DefaultSelector() as selector: + selector.register(child.stdout, selectors.EVENT_READ) + self.assertTrue(selector.select(timeout=5), 'child boundary deadline') + self.assertEqual(child.stdout.readline().strip(), 'ready') + if phase == 'sending': + self.assertTrue(self.remote.effect_committed.wait(3)) + child.kill() + child.wait(timeout=5) + finally: + if child.poll() is None: + child.kill() + child.wait(timeout=5) + child.stdout.close() + if phase == 'sending': + self.remote.release.set() + self.model.recover() + row = self.model.execute(key, {}) + self.assertEqual(row['state'], 'prepared' if phase == 'prepared' else 'outcome_unknown') + self.assertEqual(self.remote.effects, 1) + if __name__ == '__main__': unittest.main(verbosity=2) From f1174b54080cf9e5a800bdb65a19a178b3aad845 Mon Sep 17 00:00:00 2001 From: dada-yan Date: Mon, 21 Sep 2026 00:50:48 +1000 Subject: [PATCH 3/5] Propose the action dispatch contract with reproducible failure experiments Signed-off-by: dada-yan --- scripts/prototypes/actions/requirements.txt | 1 + 1 file changed, 1 insertion(+) create mode 100644 scripts/prototypes/actions/requirements.txt diff --git a/scripts/prototypes/actions/requirements.txt b/scripts/prototypes/actions/requirements.txt new file mode 100644 index 000000000..c4d28e077 --- /dev/null +++ b/scripts/prototypes/actions/requirements.txt @@ -0,0 +1 @@ +psycopg[binary]==3.2.10 From dc1145f21e62a1b371d6e662af13821e07813d8f Mon Sep 17 00:00:00 2001 From: dada-yan Date: Mon, 21 Sep 2026 00:51:18 +1000 Subject: [PATCH 4/5] Link the proposed action revision to its executable evidence Signed-off-by: dada-yan --- .../0034-an-action-is-a-declared-call.md | 7 ++ scripts/prototypes/actions/README.md | 105 ++++++++++++++++++ 2 files changed, 112 insertions(+) create mode 100644 scripts/prototypes/actions/README.md diff --git a/docs/decisions/0034-an-action-is-a-declared-call.md b/docs/decisions/0034-an-action-is-a-declared-call.md index dd5e34ff2..aec7fde1a 100644 --- a/docs/decisions/0034-an-action-is-a-declared-call.md +++ b/docs/decisions/0034-an-action-is-a-declared-call.md @@ -120,3 +120,10 @@ A run is synchronous in this cut: a person presses Run and waits for the row. Wh - **A second kind.** Email, or a tool over MCP, would share the parameters, the grants and the log and differ only in the runner. The `kind` column is the door; each kind is its own record. - **Who asks for a grant.** A base admin who wants an action has to find the deployment admin. A request queue is the shape data sources never grew either. - **Naming.** `audit_events.action` names what a person did; an action here is a declared call. The two words meet only in the ledger, where a row reads `action.run`. + +## Proposed revision · 2026-09-21 (not accepted or implemented) + +An isolated protocol experiment and the exact decisions requested are recorded in +[the experiment report](../../scripts/prototypes/actions/README.md). +This proposal does not change the accepted decisions or implementation status above. +It must be reviewed before enabling the corresponding production write/sender path. diff --git a/scripts/prototypes/actions/README.md b/scripts/prototypes/actions/README.md new file mode 100644 index 000000000..16f21b8f0 --- /dev/null +++ b/scripts/prototypes/actions/README.md @@ -0,0 +1,105 @@ +# Proposed amendment to 0034: an action attempt keeps its identity + +Status: **proposed, executable protocol experiment only**. This does not supersede +0034, add production tables/routes, or enable action sending. Refs #530. + +0034's synchronous “run, then one row” leaves two facts indistinguishable: a call +that was never sent, and a call whose remote effect happened but whose response +was lost. A client retry then risks repeating a non-idempotent action. The example +is not hypothetical: the loopback endpoint in this experiment increments its +side-effect counter, waits on a barrier, and drops the connection before replying. +The observed result is one remote effect and an unknown local outcome. Replaying +the same execution request does not increase the counter. + +## Decision requested + +Approve a stable execution request ID, revision-bound preview, and a durable +single-attempt gate before implementing the registry's sender. Prefer **no automatic +retry, no redirects, and no recovery takeover** in this first manual cut. This +narrows the redirect behavior mentioned in 0034 and must be accepted explicitly. + +The ID belongs to one explicit user operation, scoped to actor, action and KB (or +registry-test scope). Equal inputs with a new ID are a new operation; an existing +ID with different inputs/revision is a conflict. Replays reauthorize before reading +back a run. PostgreSQL's ordinary UNIQUE with a nullable scope is insufficient: +the experiment uses UNIQUE NULLS NOT DISTINCT, supported by the project's PG16. +The model has one action; production uniqueness must also include action identity. + +Every definition edit and credential rotation increments revision in the same +transaction. Preview returns revision and a non-secret rendered request without a +run or network request. The execution service renders from the same revision and +validated arguments, never a client-supplied URL/body/header set. + +## State and dispatch authority + +- `prepared`: intent persisted, no dispatch grant yet; replays only read it. +- `dispatching`: a unique token and authorization gate committed. Only that original + flow can call send after commit confirmation. A lost commit acknowledgement means + do not send, even if a later read finds dispatching. +- `not_sent`: an atomic prepared expiration/cancellation won before dispatch. +- `response_received`: observed HTTP status, independently of 2xx and body capture. +- `outcome_unknown`: dispatch may have happened; no durable observation establishes + a response. Recovery never changes this to a sendable state. + +After a short authorization/revision check and CAS transaction, no database locks +remain held across the network. Changes to permissions before this gate reject; +a revocation after it cannot recall a request. Row-lock order for action, grants, +parameters and run must be specified in the implementation. The experiment locks +one simplified definition row, **not** Utopia's production permission tables. + +Keep HTTP status, capture state (complete/truncated/read_error/absent), safe excerpt, +original dispatch token, dispatch/observation/finish timestamps and immutable safe +request snapshot separately. A late response can update unknown using the original +token; it never grants another send. The model exercises state/token/capture, not +the complete proposed audit schema or its timestamps. + +HTTP 202 is a response, not proof of business completion. A 500 may follow a remote +side effect too. A response body failure must not erase already observed headers. +A final database-write failure may leave dispatching; retry only persisting the +observation if it remains available, never the external call. + +## What was executed + +Run from `scripts/prototypes/actions` with an **isolated PostgreSQL 16 database**: + +```sh +python -m venv .venv +.venv/bin/pip install -r requirements.txt +export UTOPIA_DATABASE_URL='postgres://.../isolated_test_database' +.venv/bin/python -m unittest -v test_protocol +``` + +The model creates/drops only a random `action_model_*` schema. The HTTP server binds +loopback and the sender's target is hard-coded loopback; it cannot be repurposed as +a production managed sender. The test dependency is isolated, not a Utopia runtime +dependency. Each test shuts down its server and drops its schema. + +Linux Python 3.13 / PostgreSQL 16: **19 tests passed**. They cover concurrent duplicate +registry submissions, nullable uniqueness, KB scope, input conflict, authorization +and revision changes, insert failure, expiry/CAS competition, lost commit ack, +remote-effect/drop, 200/202/400/500, body timeout/cap, failed observation persistence, +late token observation, no redirect follow, and proxy environment isolation of this +loopback client. A separate test kills actual Python subprocesses after prepared, +after dispatch, and after the remote effect, then retries the same operation. + +Two mutations were rejected: ordinary nullable UNIQUE produces multiple dispatches; +allowing a replay to take over prepared sends an operation whose creator was lost. +Both fail their named tests, rather than merely failing to compile. + +## Not established by this model + +No production RBAC, seal/auth handling, templating, DNS pinning, direct-client policy, +production 15-second/1-MB/4-KB limits, UI, deletion retention or production migrations +were implemented or tested here. Model caps are deliberately tiny to trigger faults. +The ordinary worker has no model registration; do not copy its running-job replay +policy into action recovery. An unknown action is not a failed internal recompute. + +After approval, implement dedicated action tables and revision/grants/preview first, +then a managed sender with no retries/redirects and audited proxy policy, then UI and +real-backend authorization/E2E. Do not reuse `client_for().post()` without examining +its redirect behavior. Secrets belong in a sealed auth block; both request snapshots +and echoed response excerpts need redaction. Logs retain unknown outcomes. Rollback +first disables new dispatch, preserves runs, and cannot reverse remote effects. + +The experiment supports an **at-most-once application dispatch attempt**, not external +exactly-once execution, packet-level guarantees or arbitrary remote business semantics. From 8f3d189d46478f437df57b4a823556dfc493dcbf Mon Sep 17 00:00:00 2001 From: dada-yan Date: Mon, 21 Sep 2026 11:47:10 +1000 Subject: [PATCH 5/5] Keep action dispatch decisions in an indexed record without a second implementation. Signed-off-by: dada-yan --- .../0034-an-action-is-a-declared-call.md | 10 +- ...eeps-its-identity-and-uncertain-outcome.md | 39 +-- docs/decisions/README.md | 2 + scripts/prototypes/actions/protocol.py | 153 ---------- scripts/prototypes/actions/requirements.txt | 1 - scripts/prototypes/actions/test_protocol.py | 269 ------------------ 6 files changed, 20 insertions(+), 454 deletions(-) rename scripts/prototypes/actions/README.md => docs/decisions/0050-an-action-attempt-keeps-its-identity-and-uncertain-outcome.md (71%) delete mode 100644 scripts/prototypes/actions/protocol.py delete mode 100644 scripts/prototypes/actions/requirements.txt delete mode 100644 scripts/prototypes/actions/test_protocol.py diff --git a/docs/decisions/0034-an-action-is-a-declared-call.md b/docs/decisions/0034-an-action-is-a-declared-call.md index aec7fde1a..2b157b2a8 100644 --- a/docs/decisions/0034-an-action-is-a-declared-call.md +++ b/docs/decisions/0034-an-action-is-a-declared-call.md @@ -98,6 +98,9 @@ GET /kbs/{id}/action-runs ?action &ok &page &per A run is synchronous in this cut: a person presses Run and waits for the row. When rules fire, the same `run_action` is called from a job. +**Revision proposed 2026-09-21:** [0050](0050-an-action-attempt-keeps-its-identity-and-uncertain-outcome.md) revisits this synchronous run-then-record boundary after observing a remote effect with a lost response. It proposes durable identity and explicit uncertainty, with no automatic retry or redirect; these changes await approval and no sender is introduced. + + ## Phasing 1. **Capability.** Schema, `utopia-store::actions`, the runner in `utopia-server` (render, send, record), the routes, this record. Tests: the store's (create, grant, a viewer never sees the auth block, a run lands as a row) and the runner's against wiremock (rendering by kind, a placeholder in the host refused, an unknown argument refused, a bound enforced, a redirect into the intranet from a public host refused, the body cap, the deadline). @@ -120,10 +123,3 @@ A run is synchronous in this cut: a person presses Run and waits for the row. Wh - **A second kind.** Email, or a tool over MCP, would share the parameters, the grants and the log and differ only in the runner. The `kind` column is the door; each kind is its own record. - **Who asks for a grant.** A base admin who wants an action has to find the deployment admin. A request queue is the shape data sources never grew either. - **Naming.** `audit_events.action` names what a person did; an action here is a declared call. The two words meet only in the ledger, where a row reads `action.run`. - -## Proposed revision · 2026-09-21 (not accepted or implemented) - -An isolated protocol experiment and the exact decisions requested are recorded in -[the experiment report](../../scripts/prototypes/actions/README.md). -This proposal does not change the accepted decisions or implementation status above. -It must be reviewed before enabling the corresponding production write/sender path. diff --git a/scripts/prototypes/actions/README.md b/docs/decisions/0050-an-action-attempt-keeps-its-identity-and-uncertain-outcome.md similarity index 71% rename from scripts/prototypes/actions/README.md rename to docs/decisions/0050-an-action-attempt-keeps-its-identity-and-uncertain-outcome.md index 16f21b8f0..d543bf5ee 100644 --- a/scripts/prototypes/actions/README.md +++ b/docs/decisions/0050-an-action-attempt-keeps-its-identity-and-uncertain-outcome.md @@ -1,15 +1,12 @@ -# Proposed amendment to 0034: an action attempt keeps its identity +# 0050 · An action attempt keeps its identity and uncertain outcome -Status: **proposed, executable protocol experiment only**. This does not supersede -0034, add production tables/routes, or enable action sending. Refs #530. +- **Status**: proposed; domain contract pending review. Documentation only; no production schema, sender or routes. +- **Written**: 2026-09-21 +- **Related**: [0034](0034-an-action-is-a-declared-call.md); [PR #840](https://github.com/deeplethe/utopia/pull/840). -0034's synchronous “run, then one row” leaves two facts indistinguishable: a call -that was never sent, and a call whose remote effect happened but whose response -was lost. A client retry then risks repeating a non-idempotent action. The example -is not hypothetical: the loopback endpoint in this experiment increments its -side-effect counter, waits on a barrier, and drops the connection before replying. -The observed result is one remote effect and an unknown local outcome. Replaying -the same execution request does not increase the counter. +## Problem + +0034's synchronous run-then-record shape cannot distinguish an unsent call from a remote effect whose response was lost. Retrying the same user operation can then repeat a non-idempotent effect. Persisting intent before dispatch and preserving uncertainty makes that distinction reviewable without promising knowledge of the remote business outcome. ## Decision requested @@ -58,21 +55,9 @@ side effect too. A response body failure must not erase already observed headers A final database-write failure may leave dispatching; retry only persisting the observation if it remains available, never the external call. -## What was executed - -Run from `scripts/prototypes/actions` with an **isolated PostgreSQL 16 database**: - -```sh -python -m venv .venv -.venv/bin/pip install -r requirements.txt -export UTOPIA_DATABASE_URL='postgres://.../isolated_test_database' -.venv/bin/python -m unittest -v test_protocol -``` +## Historical investigation -The model creates/drops only a random `action_model_*` schema. The HTTP server binds -loopback and the sender's target is hard-coded loopback; it cannot be repurposed as -a production managed sender. The test dependency is isolated, not a Utopia runtime -dependency. Each test shuts down its server and drops its schema. +The investigation at `dc1145f21e62a1b371d6e662af13821e07813d8f` used Python 3.13 and PostgreSQL 16.15 on Linux with a random isolated schema and a loopback-only HTTP endpoint. The source and logs were archived outside the repository before removing the model and its dependency file. These are historical protocol-model results, not fresh Rust or production-sender tests. Linux Python 3.13 / PostgreSQL 16: **19 tests passed**. They cover concurrent duplicate registry submissions, nullable uniqueness, KB scope, input conflict, authorization @@ -103,3 +88,9 @@ first disables new dispatch, preserves runs, and cannot reverse remote effects. The experiment supports an **at-most-once application dispatch attempt**, not external exactly-once execution, packet-level guarantees or arbitrary remote business semantics. + +## Alternatives and approval boundary + +Recording only after send loses intent if the process exits. Retrying an unknown outcome can duplicate a remote effect. Transferring a prepared/dispatching attempt to recovery cannot prove the original owner did not send. Exactly-once business execution requires a remote contract this project cannot invent. The conservative first cut sacrifices automatic completion to preserve a truthful, auditable uncertainty boundary. + +Approval is requested for request identity and revision binding, the single original-flow dispatch grant, and no automatic retries, redirects or recovery takeover. Those decisions revise 0034's suggested redirect/retry behavior; they are not already accepted by moving this record. Registry, authorization and sender implementation follow only after agreement. diff --git a/docs/decisions/README.md b/docs/decisions/README.md index 6fc359926..14cd1745b 100644 --- a/docs/decisions/README.md +++ b/docs/decisions/README.md @@ -73,6 +73,7 @@ The test for writing one: if someone (including us) looks at a piece of code in | 0045 | [A time mention is resolved against its document](0045-a-time-mention-is-resolved-against-its-document.md) | Accepted · cuts 1 and 2 built (#740): a document is dated from its own text, each mention is interpreted by the model and computed by code, upload time is used nowhere · cuts 3 and 4 (grades replace the confidence gate, re-resolution and the anchor queue) not built · a time expression is a mention with its words and place; the model returns shape, anchor, offset and granularity and code computes the interval; a document carries its own date, calendars and anchors across chunks, never its upload time; unresolved mentions wait for an anchor; timelines close on resolution grade instead of confidence | | 0046 | [The app surface is MCP](0046-the-app-surface-is-mcp.md) | Decided, with the refused design kept. Asked for an app center: applications built on this knowledge, mounted, run in a sandbox, handed to a team. The answer is that the surface already exists — a personal token carries identity and scope, ten read tools serve chat and MCP from one place, `as_of` reaches every graph read, and a read returns `structuredContent` with stable ledger identities — so a coding agent builds on this base today in its own platform, its own language and its own sandbox. Refused here because the layer an app would read is being replaced under it (typed facts now come only from alignment), because 0016 closes open seams before cutting new ones, and because a catalog, an execution boundary and quotas are three other products. The shape is kept with the four gates it would have to hold (runs as the caller, egress only through a declared action, a declared clock, the existing queue) and the dead ends: a container runtime (withdrawn the day it was written — WeKnora's skills are human-written and assume a shell, and they pay for it), a Wasm component runtime (better on every axis including the determinism re-parse needs, still not built because the reason is priority), a service identity per app, an app as a saved conversation. Reopened by a named customer who needs a button inside the product, by the type layer settling, or after 0034 | | 0047 | [A rule may conclude a relation](0047-a-rule-may-conclude-a-relation.md) | Proposed 2026-09-20 · nothing built · A rule reads one entity and concludes about that same entity, so a threshold over a chain — a holding above 50% in a company that itself holds above 50% in another — cannot be written at all, and the query-time path walk that answers it produces no interval, no premises and nothing a queue can see. The conclusion becomes a **relation** between the subject and one entity reached across one declared relation, valid on the intersection of every premise interval including the join edge's. The concluded edge rejoins the pool `derive()` reads and the axiom pass runs once per round, coupling the two reasoners for the first time: 0021's cycle objection is answered with the **finiteness** argument [0030](0030-a-rule-may-read-what-a-rule-concluded.md) already put in place of acyclicity, rather than with a fixed ordering that would let a legitimate rule silently never fire. Reading a value across a hop is [0032](0032-a-rule-computes-what-it-concludes.md)'s decision, reused rather than re-decided. Negation, aggregation, a second hop and user-defined recursion stay out; the three caps in play are set by measurement in the PR that changes them | +| 0050 | [An action attempt keeps its identity and uncertain outcome](0050-an-action-attempt-keeps-its-identity-and-uncertain-outcome.md) | Proposed · durable execution identity and uncertain outcomes; no sender | | | Record | Domain | Status | |---|---|---|---| @@ -123,6 +124,7 @@ The test for writing one: if someone (including us) looks at a piece of code in | 0045 | [A time mention is resolved against its document](0045-a-time-mention-is-resolved-against-its-document.md) | time | current | | 0046 | [The app surface is MCP](0046-the-app-surface-is-mcp.md) | chat-and-mcp | current | | 0047 | [A rule may conclude a relation](0047-a-rule-may-conclude-a-relation.md) | rules | current | +| 0050 | [An action attempt keeps its identity and uncertain outcome](0050-an-action-attempt-keeps-its-identity-and-uncertain-outcome.md) | lakehouse-and-actions | proposed | The status word is whether a later record has overtaken this one; what is built is in the record's own status line. Domains are the files of [../design/](../design/README.md), where every record is dated and the status words are defined. diff --git a/scripts/prototypes/actions/protocol.py b/scripts/prototypes/actions/protocol.py deleted file mode 100644 index df6e3da5f..000000000 --- a/scripts/prototypes/actions/protocol.py +++ /dev/null @@ -1,153 +0,0 @@ -"""Isolated PostgreSQL protocol experiment, NOT Utopia's action implementation. - -No production tables, routes, credentials or generic worker integration. Targets -are a test server on loopback only. The model establishes dispatch identities and -fault boundaries; it does not establish Utopia RBAC, templates or DNS pinning. -""" -import hashlib -import http.client -import json -import os -import uuid -from contextlib import contextmanager - -import psycopg -from psycopg import sql -from psycopg.rows import dict_row - - -class Conflict(Exception): - pass - - -class Model: - def __init__(self, dsn, port): - self.dsn, self.port = dsn, port - self.schema = 'action_model_' + uuid.uuid4().hex - with psycopg.connect(dsn, autocommit=True) as c: - c.execute(sql.SQL('CREATE SCHEMA {}').format(sql.Identifier(self.schema))) - with self.connect() as c: - c.execute('''CREATE TABLE definition ( - id integer PRIMARY KEY, revision integer NOT NULL, - enabled boolean NOT NULL, granted boolean NOT NULL, - editor boolean NOT NULL); - INSERT INTO definition VALUES(1,1,true,true,true); - CREATE TABLE runs ( - id uuid PRIMARY KEY, actor text NOT NULL, scope uuid, - request_id uuid NOT NULL, fingerprint text NOT NULL, - revision integer NOT NULL, state text NOT NULL, - token uuid, status integer, capture text, excerpt text, - UNIQUE NULLS NOT DISTINCT(actor, scope, request_id));''') - - @contextmanager - def connect(self): - with psycopg.connect(self.dsn, row_factory=dict_row) as c: - c.execute(sql.SQL('SET search_path TO {}').format(sql.Identifier(self.schema))) - yield c - - def close(self): - with psycopg.connect(self.dsn, autocommit=True) as c: - c.execute(sql.SQL('DROP SCHEMA {} CASCADE').format(sql.Identifier(self.schema))) - - def definition(self, **changes): - with self.connect() as c: - for key, value in changes.items(): - assert key in ('revision', 'enabled', 'granted', 'editor') - c.execute(sql.SQL('UPDATE definition SET {}=%s').format(sql.Identifier(key)), (value,)) - - @staticmethod - def authorize(row, revision): - if not row['enabled'] or not row['granted'] or not row['editor']: - raise PermissionError('denied at dispatch boundary') - if row['revision'] != revision: - raise Conflict('preview_stale') - - def preview(self, args): - with self.connect() as c: - d = c.execute('SELECT * FROM definition WHERE id=1').fetchone() - self.authorize(d, d['revision']) - return {'revision': d['revision'], 'args': args} - - def prepare(self, request_id, args, revision=1, scope=None, actor='editor'): - fingerprint = hashlib.sha256(json.dumps([revision, args], sort_keys=True, - separators=(',', ':'), allow_nan=False).encode()).hexdigest() - with self.connect() as c: - # Replays are authorized as well. This model has one definition, not - # production role tables; it tests the cutoff, not the RBAC resolver. - d = c.execute('SELECT * FROM definition WHERE id=1 FOR SHARE').fetchone() - self.authorize(d, revision) - new = c.execute('''INSERT INTO runs(id,actor,scope,request_id,fingerprint,revision,state) - VALUES(%s,%s,%s,%s,%s,%s,'prepared') ON CONFLICT DO NOTHING RETURNING *''', - (uuid.uuid4(), actor, scope, request_id, fingerprint, revision)).fetchone() - if new: - return new, True - old = c.execute('SELECT * FROM runs WHERE actor=%s AND scope IS NOT DISTINCT FROM %s AND request_id=%s', - (actor, scope, request_id)).fetchone() - if old['fingerprint'] != fingerprint: - raise Conflict('execution_request_reused') - return old, False - - def gate(self, run_id, revision, commit_uncertain=False): - token = uuid.uuid4() - with self.connect() as c: - d = c.execute('SELECT * FROM definition WHERE id=1 FOR SHARE').fetchone() - self.authorize(d, revision) - got = c.execute("UPDATE runs SET state='dispatching',token=%s WHERE id=%s AND state='prepared' RETURNING id", - (token, run_id)).fetchone() - # Commit happened, but caller cannot know: it must NOT send. Deliberate - # fault injection after real commit models a lost commit acknowledgement. - if commit_uncertain: - raise ConnectionError('injected lost gate commit acknowledgement') - return token if got else None - - def read(self, run_id): - with self.connect() as c: - return c.execute('SELECT * FROM runs WHERE id=%s', (run_id,)).fetchone() - - def expire(self, run_id): - with self.connect() as c: - c.execute("UPDATE runs SET state='not_sent' WHERE id=%s AND state='prepared'", (run_id,)) - - def recover(self): - with self.connect() as c: - c.execute("UPDATE runs SET state='outcome_unknown' WHERE state='dispatching'") - - def observe(self, run_id, token, state, status=None, capture='absent', excerpt=''): - with self.connect() as c: - c.execute('''UPDATE runs SET state=%s,status=%s,capture=%s,excerpt=%s - WHERE id=%s AND token=%s AND state IN ('dispatching','outcome_unknown')''', - (state, status, capture, excerpt, run_id, token)) - - def send(self, run_id, token, path, fail_write=False): - # Fixed loopback destination: this is intentionally NOT a reusable sender. - # http.client neither follows redirects nor retries requests. No proxies. - conn = http.client.HTTPConnection('127.0.0.1', self.port, timeout=0.3) - status, capture, excerpt = None, 'absent', '' - try: - conn.request('POST', path, body=b'{}', headers={'Content-Type': 'application/json'}) - response = conn.getresponse() - status = response.status - try: - body = response.read(1025) - capture = 'truncated' if len(body) > 1024 else 'complete' - excerpt = body[:256].decode('utf-8', errors='replace') - except (OSError, http.client.HTTPException): - capture = 'read_error' - except (OSError, http.client.HTTPException): - pass - finally: - conn.close() - if fail_write: - raise ConnectionError('injected observation persistence failure') - self.observe(run_id, token, 'response_received' if status is not None else 'outcome_unknown', - status, capture, excerpt) - - def execute(self, request_id, args, revision=1, scope=None, path='/ok'): - row, created = self.prepare(request_id, args, revision, scope) - # Only the original creator proceeds. Returning prepared after a restart - # does not grant replay callers permission to take over and send. - if created: - token = self.gate(row['id'], revision) - if token: - self.send(row['id'], token, path) - return self.read(row['id']) diff --git a/scripts/prototypes/actions/requirements.txt b/scripts/prototypes/actions/requirements.txt deleted file mode 100644 index c4d28e077..000000000 --- a/scripts/prototypes/actions/requirements.txt +++ /dev/null @@ -1 +0,0 @@ -psycopg[binary]==3.2.10 diff --git a/scripts/prototypes/actions/test_protocol.py b/scripts/prototypes/actions/test_protocol.py deleted file mode 100644 index bf9dee5a1..000000000 --- a/scripts/prototypes/actions/test_protocol.py +++ /dev/null @@ -1,269 +0,0 @@ -import concurrent.futures -import http.server -import os -import socket -import threading -import unittest -import uuid -from protocol import Model, Conflict - - -class Remote(http.server.BaseHTTPRequestHandler): - def log_message(self, *_): - pass - - def do_POST(self): - self.rfile.read(int(self.headers.get('Content-Length', 0))) - with self.server.guard: - self.server.requests.append(self.path) - self.server.effects += 1 - if self.path == '/drop': - # A real remote effect is committed BEFORE losing the response. - self.server.effect_committed.set() - self.server.release.wait(5) - self.connection.shutdown(socket.SHUT_RDWR) - self.connection.close() - return - if self.path == '/slow-body': - self.send_response(202) - self.send_header('Content-Length', '100') - self.end_headers() - self.wfile.flush() - self.server.release.wait(2) - return - status = int(self.path[1:]) if self.path[1:].isdigit() else 200 - body = ('汉字' * 1000).encode() if self.path == '/large' else b'observed response' - self.send_response(status) - if status == 302: - self.send_header('Location', '/second-target') - self.send_header('Content-Length', str(len(body))) - self.end_headers() - self.wfile.write(body) - - -class ProtocolTests(unittest.TestCase): - def setUp(self): - self.remote = http.server.ThreadingHTTPServer(('127.0.0.1', 0), Remote) - self.remote.requests, self.remote.effects = [], 0 - self.remote.guard = threading.Lock() - self.remote.effect_committed = threading.Event() - self.remote.release = threading.Event() - self.thread = threading.Thread(target=self.remote.serve_forever, daemon=True) - self.thread.start() - self.model = Model(os.environ['UTOPIA_DATABASE_URL'], self.remote.server_port) - self.key = uuid.uuid4() - - def tearDown(self): - self.remote.release.set() - self.remote.shutdown() - self.remote.server_close() - self.thread.join() - self.model.close() - - def prepared(self): - return self.model.prepare(self.key, {'n': 1})[0] - - def test_preview_has_neither_run_nor_network_effect(self): - self.assertEqual(self.model.preview({'n': 1})['revision'], 1) - with self.model.connect() as c: - self.assertEqual(c.execute('SELECT count(*) n FROM runs').fetchone()['n'], 0) - self.assertEqual(self.remote.effects, 0) - - def test_stale_preview_and_revoked_authority_never_dispatch(self): - self.model.definition(revision=2) - with self.assertRaises(Conflict): - self.model.execute(self.key, {}, revision=1) - self.model.definition(revision=1) - row = self.prepared() - for field in ['enabled', 'granted', 'editor']: - self.model.definition(**{field: False}) - with self.assertRaises(PermissionError): - self.model.gate(row['id'], 1) - with self.assertRaises(PermissionError): - self.model.execute(self.key, {'n': 1}) - self.model.definition(**{field: True}) - self.assertEqual(self.remote.effects, 0) - - def test_concurrent_duplicate_registry_scope_has_one_run_and_one_dispatch(self): - barrier = threading.Barrier(8) - def call(_): - barrier.wait(timeout=5) - return self.model.execute(self.key, {'n': 1})['id'] - with concurrent.futures.ThreadPoolExecutor(8) as executor: - ids = list(executor.map(call, range(8))) - self.assertEqual(len(set(ids)), 1) - self.assertEqual(self.remote.effects, 1) - - def test_kb_scope_retries_and_new_explicit_operation(self): - scope = uuid.uuid4() - a = self.model.execute(self.key, {'n': 1}, scope=scope) - b = self.model.execute(self.key, {'n': 1}, scope=scope) - self.assertEqual(a['id'], b['id']) - self.model.execute(uuid.uuid4(), {'n': 1}, scope=scope) - self.assertEqual(self.remote.effects, 2) - - def test_same_identity_cannot_change_inputs(self): - self.model.execute(self.key, {'n': 1}) - with self.assertRaises(Conflict): - self.model.execute(self.key, {'n': 2}) - self.assertEqual(self.remote.effects, 1) - - def test_prepared_insert_failure_cannot_send(self): - with self.model.connect() as c: - c.execute("ALTER TABLE runs ADD CONSTRAINT injected_failure CHECK (state <> 'prepared')") - with self.assertRaises(Exception): - self.model.execute(self.key, {}) - self.assertEqual(self.remote.effects, 0) - - def test_prepared_replay_is_not_a_recovery_sender(self): - row = self.prepared() - self.assertEqual(self.model.execute(self.key, {'n': 1})['state'], 'prepared') - self.model.expire(row['id']) - self.assertIsNone(self.model.gate(row['id'], 1)) - self.assertEqual(self.remote.effects, 0) - - def test_expiration_and_dispatch_gate_are_mutually_exclusive(self): - row = self.prepared() - barrier = threading.Barrier(2) - def expire(): - barrier.wait(timeout=5) - self.model.expire(row['id']) - def gate(): - barrier.wait(timeout=5) - return self.model.gate(row['id'], 1) - with concurrent.futures.ThreadPoolExecutor(2) as ex: - a, b = ex.submit(expire), ex.submit(gate) - a.result() - token = b.result() - self.assertEqual(self.model.read(row['id'])['state'], 'dispatching' if token else 'not_sent') - self.assertIsNone(self.model.gate(row['id'], 1)) - - def test_lost_gate_commit_acknowledgement_does_not_authorize_send(self): - row = self.prepared() - with self.assertRaises(ConnectionError): - self.model.gate(row['id'], 1, commit_uncertain=True) - self.assertEqual(self.model.read(row['id'])['state'], 'dispatching') - self.model.recover() - self.assertEqual(self.model.execute(self.key, {'n': 1})['state'], 'outcome_unknown') - self.assertEqual(self.remote.effects, 0) - - def test_gate_then_crash_before_send_is_conservatively_unknown(self): - row = self.prepared() - self.model.gate(row['id'], 1) - self.model.recover() - self.assertEqual(self.model.execute(self.key, {'n': 1})['state'], 'outcome_unknown') - self.assertEqual(self.remote.effects, 0) - - def test_remote_effect_then_dropped_response_is_not_retried(self): - with concurrent.futures.ThreadPoolExecutor(1) as ex: - future = ex.submit(self.model.execute, self.key, {'n': 1}, path='/drop') - try: - self.assertTrue(self.remote.effect_committed.wait(3)) - self.assertEqual(self.remote.effects, 1) - finally: - self.remote.release.set() - row = future.result(timeout=3) - self.assertEqual(row['state'], 'outcome_unknown') - self.model.recover() - self.model.execute(self.key, {'n': 1}) - self.assertEqual(self.remote.effects, 1) - - def test_http_status_is_observation_not_business_success(self): - for status in [200, 202, 400, 500]: - row = self.model.execute(uuid.uuid4(), {}, path=f'/{status}') - self.assertEqual(row['state'], 'response_received') - self.assertEqual(row['status'], status) - self.assertEqual(row['capture'], 'complete') - self.assertEqual(self.remote.effects, 4) - - def test_body_timeout_preserves_observed_status(self): - row = self.model.execute(self.key, {}, path='/slow-body') - self.assertEqual((row['state'], row['status'], row['capture']), ('response_received', 202, 'read_error')) - self.assertEqual(self.remote.effects, 1) - - def test_body_cap_is_reported_and_excerpt_is_unicode(self): - row = self.model.execute(self.key, {}, path='/large') - self.assertEqual(row['capture'], 'truncated') - self.assertIsInstance(row['excerpt'], str) - self.assertLessEqual(len(row['excerpt']), 256) - self.assertEqual(self.remote.effects, 1) - - def test_observation_write_failure_recovery_never_resends(self): - row = self.prepared() - token = self.model.gate(row['id'], 1) - with self.assertRaises(ConnectionError): - self.model.send(row['id'], token, '/ok', fail_write=True) - self.model.recover() - self.assertEqual(self.model.execute(self.key, {'n': 1})['state'], 'outcome_unknown') - self.assertEqual(self.remote.effects, 1) - - def test_late_observation_uses_original_token_without_dispatch(self): - row = self.prepared() - token = self.model.gate(row['id'], 1) - self.model.recover() - self.model.observe(row['id'], uuid.uuid4(), 'response_received', 200) - self.assertEqual(self.model.read(row['id'])['state'], 'outcome_unknown') - self.model.observe(row['id'], token, 'response_received', 202) - self.assertEqual(self.model.execute(self.key, {'n': 1})['status'], 202) - self.assertEqual(self.remote.effects, 0) - - def test_redirect_is_an_observed_response_without_a_second_request(self): - row = self.model.execute(self.key, {}, path='/302') - self.assertEqual(row['status'], 302) - self.assertEqual(self.remote.requests, ['/302']) - - def test_proxy_environment_cannot_redirect_loopback_experiment(self): - old = os.environ.get('HTTP_PROXY') - os.environ['HTTP_PROXY'] = 'http://127.0.0.1:1' - try: - self.assertEqual(self.model.execute(self.key, {})['status'], 200) - finally: - if old is None: - os.environ.pop('HTTP_PROXY', None) - else: - os.environ['HTTP_PROXY'] = old - - def test_real_process_exit_before_and_after_dispatch_boundary(self): - import subprocess - import sys - import selectors - for phase in ['prepared', 'dispatching', 'sending']: - key = uuid.uuid4() - child_code = """ -import os,sys,time -from protocol import Model -m=Model.__new__(Model) -m.dsn=os.environ['UTOPIA_DATABASE_URL'];m.schema=sys.argv[1];m.port=int(sys.argv[2]) -row,_=m.prepare(sys.argv[3], {}) -if sys.argv[4]!='prepared': token=m.gate(row['id'],1) -print('ready',flush=True) -if sys.argv[4]=='sending': m.send(row['id'],token,'/drop') -time.sleep(30) -""" - child = subprocess.Popen([sys.executable, '-c', child_code, self.model.schema, - str(self.remote.server_port), str(key), phase], - stdout=subprocess.PIPE, text=True) - try: - with selectors.DefaultSelector() as selector: - selector.register(child.stdout, selectors.EVENT_READ) - self.assertTrue(selector.select(timeout=5), 'child boundary deadline') - self.assertEqual(child.stdout.readline().strip(), 'ready') - if phase == 'sending': - self.assertTrue(self.remote.effect_committed.wait(3)) - child.kill() - child.wait(timeout=5) - finally: - if child.poll() is None: - child.kill() - child.wait(timeout=5) - child.stdout.close() - if phase == 'sending': - self.remote.release.set() - self.model.recover() - row = self.model.execute(key, {}) - self.assertEqual(row['state'], 'prepared' if phase == 'prepared' else 'outcome_unknown') - self.assertEqual(self.remote.effects, 1) - - -if __name__ == '__main__': - unittest.main(verbosity=2)