135 lines
6.3 KiB
Python
135 lines
6.3 KiB
Python
import json
|
|
import subprocess
|
|
from pathlib import Path
|
|
|
|
import pytest
|
|
|
|
from tests.dashboard_bundle import dashboard
|
|
|
|
|
|
ROOT = Path(__file__).parents[1]
|
|
COORDINATOR = ROOT / "frontend" / "outbox-coordinator.js"
|
|
ISSUE_OUTBOX = ROOT / "frontend" / "issue-outbox.js"
|
|
AUTHORED_OUTBOX = ROOT / "frontend" / "authored-outbox.js"
|
|
|
|
|
|
def run_node(script: str):
|
|
result = subprocess.run(["node", "-e", script], check=True, capture_output=True, text=True)
|
|
return json.loads(result.stdout)
|
|
|
|
|
|
def test_two_tabs_share_one_flush_lease_for_each_outbox():
|
|
script = f"""
|
|
const createCoordinator = require({json.dumps(str(COORDINATOR))});
|
|
const createIssueOutbox = require({json.dumps(str(ISSUE_OUTBOX))});
|
|
const createAuthoredOutbox = require({json.dumps(str(AUTHORED_OUTBOX))});
|
|
const values = new Map();
|
|
const storage = {{getItem:k=>values.get(k)||null,setItem:(k,v)=>values.set(k,v),removeItem:k=>values.delete(k)}};
|
|
const held = new Set();
|
|
const locks = {{request: async (name, options, work) => {{
|
|
if (held.has(name)) return work(null);
|
|
held.add(name);
|
|
try {{ return await work({{name}}); }} finally {{ held.delete(name); }}
|
|
}}}};
|
|
let releaseIssue, releaseMessage;
|
|
const issueGate = new Promise(resolve => releaseIssue = resolve);
|
|
const messageGate = new Promise(resolve => releaseMessage = resolve);
|
|
let issueCalls = 0, messageCalls = 0;
|
|
const issueA = createIssueOutbox({{storage,getOwnerLogin:()=> 'timmy',createOperationId:()=> 'issue-1',
|
|
coordinator:createCoordinator({{storage,locks,channelFactory:null,tabId:'a'}}),fetchJson:async()=>{{issueCalls++;await issueGate;return {{number:1}};}}}});
|
|
const issueB = createIssueOutbox({{storage,getOwnerLogin:()=> 'timmy',
|
|
coordinator:createCoordinator({{storage,locks,channelFactory:null,tabId:'b'}}),fetchJson:async()=>{{issueCalls++;return {{number:1}};}}}});
|
|
issueA.enqueue({{repository:'o/r',title:'One'}});
|
|
const messageA = createAuthoredOutbox({{storage,getOwnerLogin:()=> 'timmy',
|
|
coordinator:createCoordinator({{storage,locks,channelFactory:null,tabId:'a'}}),fetchJson:async()=>{{messageCalls++;await messageGate;return {{id:1}};}}}});
|
|
const messageB = createAuthoredOutbox({{storage,getOwnerLogin:()=> 'timmy',
|
|
coordinator:createCoordinator({{storage,locks,channelFactory:null,tabId:'b'}}),fetchJson:async()=>{{messageCalls++;return {{id:1}};}}}});
|
|
messageA.enqueue({{kind:'issue-comment',repository:'o/r',number:1,body:'Reply',operationId:'message-1'}});
|
|
(async()=>{{
|
|
const issueFirst=issueA.flush('timmy'); const issueSecond=issueB.retry('issue-1', 'timmy');
|
|
const messageFirst=messageA.flush('timmy'); const messageSecond=messageB.retry('message-1', 'timmy');
|
|
await Promise.resolve();
|
|
releaseIssue(); releaseMessage();
|
|
const results=await Promise.all([issueFirst,issueSecond,messageFirst,messageSecond]);
|
|
process.stdout.write(JSON.stringify({{issueCalls,messageCalls,results,issueRemaining:issueA.list(),messageRemaining:messageA.list()}}));
|
|
}})();
|
|
"""
|
|
output = run_node(script)
|
|
|
|
assert output["issueCalls"] == 1
|
|
assert output["messageCalls"] == 1
|
|
assert output["issueRemaining"] == []
|
|
assert output["messageRemaining"] == []
|
|
assert sum(result.get("lease_skipped", False) for result in output["results"]) == 2
|
|
|
|
|
|
def test_fallback_lease_expires_and_stale_owner_cannot_release_successor():
|
|
script = f"""
|
|
const createCoordinator = require({json.dumps(str(COORDINATOR))});
|
|
const values = new Map();
|
|
const storage = {{getItem:k=>values.get(k)||null,setItem:(k,v)=>values.set(k,v),removeItem:k=>values.delete(k)}};
|
|
let clock=100;
|
|
let releaseFirst;
|
|
const firstGate=new Promise(resolve=>releaseFirst=resolve);
|
|
const first=createCoordinator({{storage,locks:null,channelFactory:null,tabId:'first',now:()=>clock,leaseMs:50}});
|
|
const second=createCoordinator({{storage,locks:null,channelFactory:null,tabId:'second',now:()=>clock,leaseMs:50}});
|
|
(async()=>{{
|
|
const running=first.runExclusive('issue', async()=>{{await firstGate;return 'first-done';}});
|
|
await Promise.resolve();
|
|
const blocked=await second.runExclusive('issue', async()=> 'too-early');
|
|
clock=151;
|
|
const takeover=await second.runExclusive('issue', async()=> 'taken-over');
|
|
releaseFirst();
|
|
const original=await running;
|
|
const after=await second.runExclusive('issue', async()=> 'after-release');
|
|
process.stdout.write(JSON.stringify({{blocked,takeover,original,after,lease:storage.getItem('stackchain.outbox-lease.issue.v1')}}));
|
|
}})();
|
|
"""
|
|
output = run_node(script)
|
|
|
|
assert output["blocked"]["lease_skipped"] is True
|
|
assert output["takeover"] == "taken-over"
|
|
assert output["original"] == "first-done"
|
|
assert output["after"] == "after-release"
|
|
assert output["lease"] is None
|
|
|
|
|
|
def test_queue_change_notifications_cross_tabs_and_can_unsubscribe():
|
|
script = f"""
|
|
const createCoordinator = require({json.dumps(str(COORDINATOR))});
|
|
const values = new Map();
|
|
const storage = {{getItem:k=>values.get(k)||null,setItem:(k,v)=>values.set(k,v),removeItem:k=>values.delete(k)}};
|
|
const channels=[];
|
|
function channelFactory() {{
|
|
const channel={{onmessage:null,postMessage(message){{channels.filter(x=>x!==channel).forEach(x=>x.onmessage?.({{data:message}}));}},close(){{}}}};
|
|
channels.push(channel); return channel;
|
|
}}
|
|
const first=createCoordinator({{storage,channelFactory,tabId:'first'}});
|
|
const second=createCoordinator({{storage,channelFactory,tabId:'second'}});
|
|
const seen=[];
|
|
const unsubscribe=second.subscribe(change=>seen.push(change));
|
|
first.notify('issue');
|
|
first.notify('authored');
|
|
unsubscribe();
|
|
first.notify('issue');
|
|
process.stdout.write(JSON.stringify({{seen,signal:JSON.parse(storage.getItem('stackchain.outbox-change.v1'))}}));
|
|
"""
|
|
output = run_node(script)
|
|
|
|
assert [change["queue"] for change in output["seen"]] == ["issue", "authored"]
|
|
assert output["signal"]["queue"] == "issue"
|
|
assert output["signal"]["tabId"] == "first"
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_dashboard_wires_cross_tab_draft_refresh_and_precaches_coordinator():
|
|
html = await dashboard()
|
|
worker = (ROOT / "frontend" / "service-worker.js").read_text()
|
|
|
|
assert '<script src="static/outbox-coordinator.js"></script>' in html
|
|
assert "const outboxCoordinator = createOutboxCoordinator" in html
|
|
assert "coordinator: outboxCoordinator" in html
|
|
assert "outboxCoordinator.subscribe" in html
|
|
assert "refreshMyWorkView()" in html
|
|
assert "BASE + 'static/outbox-coordinator.js'" in worker
|