| 1 | #!/usr/bin/env python3 |
| 2 | """Compact a shared section between two halves of a twelve-client editing run.""" |
| 3 | import argparse |
| 4 | import base64 |
| 5 | import json |
| 6 | import os |
| 7 | from pathlib import Path |
| 8 | import signal |
| 9 | import subprocess |
| 10 | import sys |
| 11 | import time |
| 12 | import xml.etree.ElementTree as ET |
| 13 | |
| 14 | from native_runner import ROOT, windows |
| 15 | import linux_vm |
| 16 | from verify_smb_overlap import verify |
| 17 | |
| 18 | |
| 19 | def maintenance_locks(events): |
| 20 | pending, files, peers = {}, {}, {} |
| 21 | phase, attempts = None, [] |
| 22 | for event in events: |
| 23 | assert not event.get('trace_error') and not event.get('encrypted') |
| 24 | phase = event.get('control', {}).get('phase', phase) |
| 25 | connection = event.get('connection') |
| 26 | if event.get('opened'): peers[connection] = event['peer'][0] |
| 27 | if event.get('closed'): |
| 28 | files = {key: value for key, value in files.items() if key[0] != connection} |
| 29 | if 'command' not in event: continue |
| 30 | key = connection, event['message'] |
| 31 | if event['direction'] == 'request': |
| 32 | pending[key] = event, phase |
| 33 | continue |
| 34 | if event['status'] == '0x103': continue |
| 35 | pair = pending.pop(key, None) |
| 36 | if pair is None: continue |
| 37 | request, issued = pair |
| 38 | command = request['command'] |
| 39 | if command == 5 and event['status'] == '0x0': |
| 40 | files[connection, event['file_id']] = request['path'].replace('\\', '/').lower() |
| 41 | elif command == 6: |
| 42 | files.pop((connection, request['file_id']), None) |
| 43 | elif command == 10 and peers[connection].startswith('192.168.77.'): |
| 44 | path = files.get((connection, request['file_id']), '') |
| 45 | if path != 'm6-collaboration/synthetic.one': continue |
| 46 | for offset, length, flags in request['locks']: |
| 47 | if (offset, length) == (0xffffeffc, 4096) and flags & 3 == 2: |
| 48 | attempts.append({'phase': issued, 'status': event['status'], 'connection': connection, |
| 49 | 'message': event['message'], 'flags': flags}) |
| 50 | assert any(a['phase'] == 'maintenance-held' and a['status'] in ('0xc0000054', '0xc0000055') for a in attempts), 'No section maintenance conflict observed while the reader held its guard' |
| 51 | assert any(a['phase'] == 'maintenance-released' and a['status'] == '0x0' for a in attempts), 'No section maintenance guard acquired after reader release' |
| 52 | return attempts |
| 53 | |
| 54 | |
| 55 | def run(output, server, fixture, operations, seed): |
| 56 | output = output.resolve() |
| 57 | assert not output.exists(), 'Choose a new output directory' |
| 58 | output.parent.mkdir(parents=True, exist_ok=True) |
| 59 | guardian = None |
| 60 | with output.with_suffix('.log').open('x') as log: |
| 61 | process = subprocess.Popen([sys.executable, ROOT / 'tools/native_collaboration.py', output, |
| 62 | '--linux', server, '--stress-clients', '4', '--rust-writers', '4', '--rust-readers', '4', |
| 63 | '--stress-operations', str(operations), '--seed', str(seed), '--edit', '--embedded-smb', |
| 64 | '--maintenance', '--fixture', fixture], stdout=log, stderr=subprocess.STDOUT) |
| 65 | |
| 66 | def wait_for(predicate, timeout, message): |
| 67 | deadline = time.monotonic() + timeout |
| 68 | while not predicate(): |
| 69 | if process.poll() is not None: raise RuntimeError(f'Workload exited with {process.returncode}: {message}') |
| 70 | if guardian is not None and guardian.poll() is not None and guardian.returncode != 0: |
| 71 | raise RuntimeError('The reader guardian failed') |
| 72 | if time.monotonic() > deadline: raise TimeoutError(message) |
| 73 | time.sleep(.2) |
| 74 | |
| 75 | def ssh(command): |
| 76 | result = linux_vm.run_ssh(server, command, timeout=30) |
| 77 | with (output / 'maintenance-server.jsonl').open('a') as stream: |
| 78 | stream.write(json.dumps({'command': command, 'exit': result.returncode, 'stdout': result.stdout, 'stderr': result.stderr}) + '\n') |
| 79 | result.check_returncode() |
| 80 | return result.stdout |
| 81 | |
| 82 | def phase(name): |
| 83 | ssh('printf \'{"phase":"' + name + '"}\' > /tmp/smb-control.json') |
| 84 | wait_for(lambda: name in ssh(f'grep -F \'"control": {{"phase": "{name}"}}\' /tmp/smb-trace.jsonl || true'), |
| 85 | 10, 'Proxy did not acknowledge the maintenance phase') |
| 86 | |
| 87 | def stat(): |
| 88 | inode, size = ssh("stat -c '%i %s' /srv/agent/m6-collaboration/synthetic.one").split() |
| 89 | return {'inode': int(inode), 'size': int(size)} |
| 90 | |
| 91 | def ui(name, script): |
| 92 | (output / f'{name}.ahk').write_text(script) |
| 93 | result = windows.do_exec(script, target=target, timeout_ms=60000, shot_delay_ms=500) |
| 94 | screenshot = result.pop('png_b64', None) |
| 95 | if screenshot: (output / f'{name}.png').write_bytes(base64.b64decode(screenshot)) |
| 96 | (output / f'{name}.json').write_text(json.dumps(result, indent=2)) |
| 97 | if result.get('error') or result.get('exit') != 0: raise RuntimeError(str(result)) |
| 98 | |
| 99 | try: |
| 100 | wait_for(lambda: (output / 'maintenance-ready').exists(), 900, 'Native clients did not reach the initial checkpoint') |
| 101 | (output / 'maintenance-start').touch() |
| 102 | shared, rust = output / 'mount/m6-collaboration', output / 'rust' |
| 103 | wait_for(lambda: all((shared / f'maintenance-paused-n{i}').exists() and (rust / f'paused-w{i}').exists() for i in range(4)), |
| 104 | 300, 'Writers did not reach the maintenance barrier') |
| 105 | (rust / 'pause').touch() |
| 106 | wait_for(lambda: all((rust / f'paused-r{i}').exists() for i in range(4)), 60, 'Readers did not pause') |
| 107 | target = json.loads((output / 'n0/machine.json').read_text())['name'] |
| 108 | page = ET.parse(sorted((output / 'n0').glob('snapshot-*.xml'))[-1]).getroot().attrib['ID'] |
| 109 | ui('maintenance-options', f'''if A_ScreenWidth != 800 || A_ScreenHeight != 600 |
| 110 | throw Error("The maintenance controller requires an 800 by 600 desktop.") |
| 111 | ComObject("OneNote.Application").NavigateTo("{page}", "", false) |
| 112 | WinWait("ahk_exe ONENOTE.EXE", , 10) |
| 113 | WinActivate("ahk_exe ONENOTE.EXE") |
| 114 | WinWaitActive("ahk_exe ONENOTE.EXE", , 10) |
| 115 | Send("!ft") |
| 116 | WinWait("OneNote Options", , 10) |
| 117 | WinActivate("OneNote Options") |
| 118 | WinWaitActive("OneNote Options", , 10) |
| 119 | Sleep(1000) |
| 120 | CoordMode("Mouse", "Screen") |
| 121 | Click(74, 134) |
| 122 | Sleep(1000) |
| 123 | ''') |
| 124 | hold = output / 'guardian' |
| 125 | hold.mkdir() |
| 126 | config = json.loads((output / 'linux.json').read_text()) |
| 127 | with (hold / 'run.log').open('w') as guardian_log: |
| 128 | guardian = subprocess.Popen(['cargo', 'test', '-p', 'notebook', '--features', 'smb', '--lib', 'live_reader_hold', '--', '--ignored', '--nocapture'], |
| 129 | cwd=ROOT, stdout=guardian_log, stderr=subprocess.STDOUT, |
| 130 | env={**os.environ, 'ONESTORE_SMB_LAB': f'127.0.0.1:{config["samba_port"]}', |
| 131 | 'ONESTORE_SMB_PATH': 'm6-collaboration/synthetic.one', 'ONESTORE_SMB_HOLD': str(hold)}) |
| 132 | wait_for(lambda: (hold / 'ready').exists(), 60, 'Reader guard was not acquired') |
| 133 | before = stat() |
| 134 | phase('maintenance-held') |
| 135 | ui('maintenance-denied', '''CoordMode("Mouse", "Screen") |
| 136 | WinActivate("OneNote Options") |
| 137 | WinWaitActive("OneNote Options", , 10) |
| 138 | Click(254, 440) |
| 139 | if !WinWait("Microsoft OneNote ahk_class #32770", , 30) |
| 140 | throw Error("The maintenance conflict dialog did not appear.") |
| 141 | FileAppend(WinGetText("Microsoft OneNote ahk_class #32770"), "*") |
| 142 | ''') |
| 143 | held = stat() |
| 144 | assert held == before, 'Section changed while the reader held its maintenance exclusion guard' |
| 145 | (hold / 'release').touch() |
| 146 | wait_for(lambda: guardian.poll() is not None, 30, 'Reader guardian did not release') |
| 147 | assert guardian.returncode == 0 |
| 148 | assert json.loads((hold / 'released.json').read_text())['accepted'] > 0 |
| 149 | phase('maintenance-released') |
| 150 | ui('maintenance-dismiss', '''WinActivate("Microsoft OneNote ahk_class #32770") |
| 151 | ControlClick("Button1", "Microsoft OneNote ahk_class #32770") |
| 152 | if !WinWaitClose("Microsoft OneNote ahk_class #32770", , 10) |
| 153 | throw Error("The maintenance conflict dialog did not close.") |
| 154 | ''') |
| 155 | replacements = [] |
| 156 | for attempt in range(5): |
| 157 | ui(f'maintenance-optimize-{attempt}', '''CoordMode("Mouse", "Screen") |
| 158 | WinActivate("OneNote Options") |
| 159 | WinWaitActive("OneNote Options", , 10) |
| 160 | Click(254, 440) |
| 161 | Sleep(3000) |
| 162 | if WinExist("Microsoft OneNote ahk_class #32770") { |
| 163 | FileAppend(WinGetText("Microsoft OneNote ahk_class #32770"), "*") |
| 164 | ControlClick("Button1", "Microsoft OneNote ahk_class #32770") |
| 165 | if !WinWaitClose("Microsoft OneNote ahk_class #32770", , 10) |
| 166 | throw Error("The maintenance conflict dialog did not close.") |
| 167 | } |
| 168 | ''') |
| 169 | replacements.append(stat()) |
| 170 | if replacements[-1]['inode'] != before['inode'] and replacements[-1]['size'] < before['size']: break |
| 171 | time.sleep(2) |
| 172 | (output / 'maintenance-stat.json').write_text(json.dumps({'before': before, 'held': held, 'attempts': replacements}, indent=2)) |
| 173 | assert replacements[-1]['inode'] != before['inode'] and replacements[-1]['size'] < before['size'], 'Native optimization did not replace and shrink the section' |
| 174 | ui('maintenance-options-close', '''WinActivate("OneNote Options") |
| 175 | CoordMode("Mouse", "Screen") |
| 176 | Click(663, 533) |
| 177 | WinWaitClose("OneNote Options", , 10) |
| 178 | ''') |
| 179 | phase('maintenance-resumed') |
| 180 | (rust / 'resume').touch() |
| 181 | (shared / 'maintenance-resume').touch() |
| 182 | wait_for(lambda: process.poll() is not None, 600, 'Post-maintenance workload did not complete') |
| 183 | assert process.returncode == 0, 'Mixed workload failed after maintenance' |
| 184 | events = [json.loads(line) for line in (output / 'smb-trace.jsonl').read_text().splitlines()] |
| 185 | result = {'maintenance_locks': maintenance_locks(events), 'resumed_overlap': verify(events, phase='maintenance-resumed')} |
| 186 | (output / 'maintenance-verification.json').write_text(json.dumps(result, indent=2)) |
| 187 | finally: |
| 188 | if guardian is not None and guardian.poll() is None: |
| 189 | (output / 'guardian/release').touch() |
| 190 | try: guardian.wait(timeout=30) |
| 191 | except subprocess.TimeoutExpired: |
| 192 | guardian.terminate() |
| 193 | guardian.wait(timeout=30) |
| 194 | if process.poll() is None: |
| 195 | process.terminate() |
| 196 | process.wait(timeout=180) |
| 197 | |
| 198 | |
| 199 | if __name__ == '__main__': |
| 200 | parser = argparse.ArgumentParser(description=__doc__) |
| 201 | parser.add_argument('output', type=Path) |
| 202 | parser.add_argument('--linux', required=True) |
| 203 | parser.add_argument('--fixture', type=Path, required=True) |
| 204 | parser.add_argument('--operations', type=int, default=40) |
| 205 | parser.add_argument('--seed', type=int, default=908) |
| 206 | args = parser.parse_args() |
| 207 | if args.operations < 4 or args.operations % 2: parser.error('Use an even operation count of at least four.') |
| 208 | def interrupted(_signal, _frame): raise KeyboardInterrupt |
| 209 | signal.signal(signal.SIGTERM, interrupted) |
| 210 | run(args.output, args.linux, args.fixture.resolve(), args.operations, args.seed) |