diff --git a/corpus/skills/cat-mode/scripts/update_fleet.sh b/corpus/skills/cat-mode/scripts/update_fleet.sh index 03cd24147..732359f43 100755 --- a/corpus/skills/cat-mode/scripts/update_fleet.sh +++ b/corpus/skills/cat-mode/scripts/update_fleet.sh @@ -1,6 +1,29 @@ #!/bin/bash set -uo pipefail +usage() { + cat <<'USAGE' +update_fleet.sh -- put every machine on one Invoker release and the current +catstack. + + update_fleet.sh [--version ] [--hosts ] [--skip-invoker] + [--skip-catstack] [--with-app] [--dry-run] + + --version release tag to install (default: newest daily-* release) + --hosts subset of remoteTargets ids (default: all of them) + --skip-invoker leave the Invoker CLI where it is + --skip-catstack leave the catstack checkout where it is + --with-app also replace /Applications/Invoker.app on the Mac. This + quits a running Invoker, the live owner on that machine. An + interrupted replace parks the live bundle; the run puts it + back, and so does the next run if it was killed outright. + --dry-run resolve versions and print the table; change nothing. + +Every host gets one row. A row that could not be checked says so; it never +reads as ok. Exit is non-zero if any row failed. +USAGE +} + REPO="${INVOKER_RELEASE_REPO:-Neko-Catpital-Labs/Invoker}" CONFIG="${INVOKER_CONFIG:-$HOME/.invoker/config.json}" APP_DIR="${INVOKER_APP_DIR:-/Applications}" @@ -17,22 +40,6 @@ FAILED=0 cleanup() { rm -rf "$WORK_DIR"; } trap cleanup EXIT -usage() { - cat <<'USAGE' -Put every machine on one Invoker release and the current catstack. - - update_fleet.sh [--version ] [--hosts ] [--skip-invoker] - [--skip-catstack] [--with-app] [--dry-run] - ---version release tag to install (default: newest daily-* release) ---hosts subset of remoteTargets ids (default: all of them) ---skip-invoker ---skip-catstack ---with-app also replace /Applications/Invoker.app on the Mac ---dry-run check every host and print the table; change nothing -USAGE -} - while [ "$#" -gt 0 ]; do case "$1" in --version) VERSION="${2:?--version needs a tag}"; shift 2 ;; @@ -156,18 +163,44 @@ local_invoker() { row ok local "invoker $before -> $after" "$link" } +app_version() { + defaults read "$APP_DIR/Invoker.app/Contents/Info.plist" CFBundleShortVersionString 2>/dev/null || echo none +} + +parked_app() { + local candidate + [ -d "$APP_DIR/Invoker.app" ] && return 0 + for candidate in "$APP_DIR"/Invoker.app.replacing.*; do + [ -d "$candidate" ] || continue + printf '%s' "$candidate" + return 0 + done + return 0 +} + restore_parked_app() { local backup="$1" - [ -d "$backup" ] || return 0 + [ -n "$backup" ] && [ -d "$backup" ] || return 2 rm -rf "$APP_DIR/Invoker.app" - mv "$backup" "$APP_DIR/Invoker.app" + mv "$backup" "$APP_DIR/Invoker.app" || return 1 + return 0 } local_app() { - local dmg mount app before after backup + local dmg mount app before after backup parked rc [ "$(uname -s)" = "Darwin" ] || { row skip local "app: not macOS" ""; return 0; } - before="$(defaults read "$APP_DIR/Invoker.app/Contents/Info.plist" CFBundleShortVersionString 2>/dev/null || echo none)" - if [ "$DRY_RUN" = 1 ]; then row ok local "app $before -> $RELEASE_VERSION (dry-run)" ""; return 0; fi + parked="$(parked_app)" + if [ "$DRY_RUN" = 1 ]; then + if [ -n "$parked" ]; then + row warn local "app: an interrupted run left $parked and no $APP_DIR/Invoker.app; a real run puts it back first (dry-run)" "$parked" + return 0 + fi + row ok local "app $(app_version) -> $RELEASE_VERSION (dry-run)" ""; return 0 + fi + if [ -n "$parked" ] && ! restore_parked_app "$parked"; then + row fail local "app: $APP_DIR/Invoker.app is missing and $parked could not be moved back" "$parked"; return 1 + fi + before="$(app_version)" case "$(uname -m)" in arm64) dmg="Invoker-$RELEASE_VERSION-arm64.dmg" ;; *) dmg="Invoker-$RELEASE_VERSION-x64.dmg" ;; @@ -188,25 +221,40 @@ local_app() { hdiutil detach "$mount" >/dev/null 2>&1 row fail local "app: could not move $APP_DIR/Invoker.app aside; nothing replaced" ""; return 1 fi + trap 'restore_parked_app "$APP_DIR/Invoker.app.replacing.$$"; exit 130' INT TERM HUP if ! cp -R "$app" "$APP_DIR/"; then hdiutil detach "$mount" >/dev/null 2>&1 rm -rf "$APP_DIR/Invoker.app" - if ! restore_parked_app "$backup"; then + restore_parked_app "$backup"; rc="$?" + trap - INT TERM HUP + if [ "$rc" = 1 ]; then row fail local "app: copy failed and the only bundle left is $backup" "$backup"; return 1 fi + if [ "$rc" = 2 ]; then + row fail local "app: copy failed and $APP_DIR has no Invoker.app to put back" ""; return 1 + fi row fail local "app $before unchanged: could not copy the new bundle into $APP_DIR" ""; return 1 fi hdiutil detach "$mount" >/dev/null 2>&1 - after="$(defaults read "$APP_DIR/Invoker.app/Contents/Info.plist" CFBundleShortVersionString 2>/dev/null || echo unknown)" + after="$(app_version)" if [ "$after" != "$RELEASE_VERSION" ]; then rm -rf "$APP_DIR/Invoker.app" - if ! restore_parked_app "$backup"; then + restore_parked_app "$backup"; rc="$?" + trap - INT TERM HUP + if [ "$rc" = 1 ]; then row fail local "app $before -> $after (wanted $RELEASE_VERSION); the only bundle left is $backup" "$backup"; return 1 fi + if [ "$rc" = 2 ]; then + row fail local "app $before -> $after (wanted $RELEASE_VERSION); $APP_DIR has no Invoker.app to put back" ""; return 1 + fi row fail local "app $before unchanged: copied bundle read $after (wanted $RELEASE_VERSION)" ""; return 1 fi rm -rf "$APP_DIR/Invoker.app.old" - if [ -d "$backup" ]; then mv "$backup" "$APP_DIR/Invoker.app.old"; fi + if [ -d "$backup" ] && ! mv "$backup" "$APP_DIR/Invoker.app.old"; then + trap - INT TERM HUP + row fail local "app $before -> $after, but the previous bundle is still parked at $backup" "$backup"; return 1 + fi + trap - INT TERM HUP row ok local "app $before -> $after" "relaunch it to restore the owner" } diff --git a/product/skills/event-wait/scripts/wait_event.py b/product/skills/event-wait/scripts/wait_event.py index 7b9f38925..f7dae8b2c 100644 --- a/product/skills/event-wait/scripts/wait_event.py +++ b/product/skills/event-wait/scripts/wait_event.py @@ -206,14 +206,26 @@ def parse_framing(raw, where: str) -> dict: raise SpecError("spec_range", f"{where}.kind must be lines or length_prefix") +def _body_path(raw, where: str) -> list[str]: + """Where the body sits inside a record; empty means the record is the body. + + Omitting the key and writing an explicit [] are the same request, and + references/framed-socket-source.md documents [] for a source whose records + are already the body. dig() treats an empty path that way, so the only + thing that ever rejected it was this validator. + """ + if raw is None or raw == []: + return [] + return _require_path_list(raw, where) + + def parse_envelope(raw, where: str) -> dict: if raw is None: return {"match_fields": {}, "body_path": []} envelope = _require_dict(raw, where) _reject_unknown(envelope, ("match_fields", "body_path"), where) match_fields = _require_match_fields(envelope.get("match_fields", {}), f"{where}.match_fields") - body_raw = envelope.get("body_path") - body_path = [] if body_raw is None else _require_path_list(body_raw, f"{where}.body_path") + body_path = _body_path(envelope.get("body_path"), f"{where}.body_path") return {"match_fields": match_fields, "body_path": body_path} @@ -250,8 +262,7 @@ def parse_snapshot(raw, where: str) -> dict | None: error_match = _require_match_fields( snapshot.get("error_match_fields", {}), f"{where}.error_match_fields" ) - body_raw = snapshot.get("body_path") - body_path = [] if body_raw is None else _require_path_list(body_raw, f"{where}.body_path") + body_path = _body_path(snapshot.get("body_path"), f"{where}.body_path") return { "request": dict(request), "request_id_field": request_id_field, @@ -773,6 +784,14 @@ def send_snapshot(self) -> None: self.snapshot_open = True def is_duplicate(self, body: dict) -> bool: + """Has this exact event already been seen? Only ever asked of our own subject. + + Event identity is only unique within a subject: a shared channel can + carry another subject's event under an identifier ours will reuse + later. Recording foreign identifiers here would let a neighbour's + event mark this wait's own completion as already seen, so the wait + would sit out a job that had already finished. + """ path = self.spec["match"]["event_id_path"] if path is None: return False @@ -788,14 +807,19 @@ def is_duplicate(self, body: dict) -> bool: self.seen_event_ids[key] = True return False - def terminal_status(self, body) -> str | None: - """The terminal status this body reports for our subject, if any.""" + def is_subject(self, body) -> bool: + """Does this body report on the exact subject this wait owns?""" if not isinstance(body, dict): - return None + return False match = self.spec["match"] subject = dig(body, match["subject_path"]) - if subject is MISSING or subject != match["subject"]: + return subject is not MISSING and subject == match["subject"] + + def terminal_status(self, body) -> str | None: + """The terminal status this body reports for our subject, if any.""" + if not self.is_subject(body): return None + match = self.spec["match"] status = dig(body, match["status_path"]) if status is MISSING or not isinstance(status, str): return None @@ -864,6 +888,8 @@ def handle(self, record: dict) -> dict | None: ) if kind == "snapshot_response": self.snapshot_open = False + if not self.is_subject(body): + return None if self.is_duplicate(body): return None status = self.terminal_status(body) @@ -874,6 +900,8 @@ def handle(self, record: dict) -> dict | None: return None if self.snapshot_open: self.events_during_snapshot += 1 + if not self.is_subject(body): + return None if self.is_duplicate(body): return None status = self.terminal_status(body) @@ -894,6 +922,13 @@ def excerpt(self, body) -> tuple[str, bool]: return raw[:limit].decode("utf-8", "ignore"), True def deliver_wake(self, status: str) -> tuple[bool, str]: + """Run the wake command once, keeping its output off the record stream. + + stdout carries the armed record and the receipt, and a caller parses + those lines as JSON. A wake command that prints anything would be read + as a malformed record, so its output goes to stderr, where it stays + readable when a wake has to be debugged. + """ wake = self.spec["wake"] if wake["mode"] == "none": return False, "session_wake_unsupported" @@ -902,7 +937,10 @@ def deliver_wake(self, status: str) -> tuple[bool, str]: argv.append(self.spec["receipt_path"]) try: completed = subprocess.run( - argv, timeout=wake["timeout_seconds"], env=self.wake_env(status) + argv, + timeout=wake["timeout_seconds"], + stdout=sys.stderr, + env=self.wake_env(status), ) except subprocess.TimeoutExpired: print( diff --git a/product/skills/event-wait/tests/test_wait_event.py b/product/skills/event-wait/tests/test_wait_event.py index 87b45b005..9457a72cd 100644 --- a/product/skills/event-wait/tests/test_wait_event.py +++ b/product/skills/event-wait/tests/test_wait_event.py @@ -398,6 +398,29 @@ def test_duplicate_events_are_skipped_by_event_identity(self): self.assertEqual(receipt["outcome"], "matched") self.assertEqual(receipt["duplicate_events_skipped"], 2) + def test_another_subjects_event_id_never_masks_our_completion(self): + """A neighbour on the same channel may reuse an id we have yet to send. + + Both events carry eventId seq-1, but only the second one is ours. If + the foreign event were recorded as seen, ours would read as a repeat + and the wait would sit out a job that had already finished. + """ + producer = self.producer() + spec = self.socket_spec("shared-ids", "wf-mine", producer.path, deadline_seconds=6.0) + wait = self.start(spec) + wait.read_armed() + + producer.emit({"workflowId": "wf-theirs", "status": "completed", "eventId": "seq-1"}) + time.sleep(0.3) + producer.emit({"workflowId": "wf-mine", "status": "completed", "eventId": "seq-1"}) + + receipt = wait.finish() + self.assertEqual(receipt["outcome"], "matched") + self.assertEqual(receipt["subject"], "wf-mine") + self.assertEqual(receipt["status"], "completed") + self.assertEqual(receipt["duplicate_events_skipped"], 0) + self.assertEqual(receipt["exit_code"], 0) + def test_repeated_terminal_event_still_writes_one_receipt(self): producer = self.producer() spec = self.socket_spec("one-receipt", "wf-1", producer.path) @@ -647,6 +670,27 @@ def test_a_clean_disconnect_before_any_match_is_a_source_error(self): self.assertFalse(receipt["event_received"]) self.assertFalse(receipt["wake_delivered"]) + def test_a_bad_frame_later_in_a_read_does_not_bury_an_earlier_match(self): + """Same read, but the trailing frame is unparsable rather than oversize. + + The receipt must also stay clean: the junk frame carries a secret, and + nothing from an undecodable frame belongs in an error the caller reads. + """ + producer = self.producer() + spec = self.socket_spec("late-garbage", "wf-1", producer.path, deadline_seconds=8.0) + wait = self.start(spec) + wait.read_armed() + producer.send_raw( + pub({"workflowId": "wf-1", "status": "completed", "eventId": "g1"}) + + frame(b"{not json SUPERSECRET") + ) + + receipt = wait.finish() + self.assertEqual(receipt["outcome"], "matched") + self.assertEqual(receipt["status"], "completed") + self.assertEqual(receipt["exit_code"], wait_event.EXIT_MATCHED) + self.assertNotIn("SUPERSECRET", json.dumps(receipt)) + def test_a_missing_socket_reports_an_unarmed_source(self): spec = self.socket_spec( "no-socket", "wf-1", os.path.join(self.dir, "absent.sock"), deadline_seconds=5.0 @@ -757,6 +801,35 @@ def test_no_declared_wake_reports_the_mechanism_as_unsupported(self): self.assertEqual(receipt["wake_status"], "session_wake_unsupported") self.assertEqual(receipt["exit_code"], 0) + def test_a_chatty_wake_never_writes_onto_the_record_stream(self): + """A wake that prints must not corrupt the JSON the caller parses. + + stdout carries the armed record and the receipt. If the wake command's + own output landed there, the caller would read it as a malformed + record instead of a receipt. + """ + producer = self.producer() + spec = self.socket_spec( + "wake-chatty", + "wf-1", + producer.path, + wake={ + "mode": "command", + "owner": "session-owner-1", + "argv": [sys.executable, "-c", "print('WAKE-NOISE')"], + "ready_argv": [sys.executable, "-c", "print('PROBE-NOISE')"], + }, + ) + wait = self.start(spec) + self.assertTrue(wait.read_armed()["callback_ready"]) + producer.emit({"workflowId": "wf-1", "status": "completed", "eventId": "w9"}) + + receipt = wait.finish() + self.assertEqual(receipt["record"], "receipt") + self.assertTrue(receipt["wake_delivered"]) + self.assertEqual(receipt["wake_status"], "delivered") + self.assertEqual(wait.process.returncode, wait_event.EXIT_MATCHED) + def test_a_delivered_wake_runs_once_and_sees_only_spec_derived_values(self): producer = self.producer() marker = os.path.join(self.dir, "woke.json") @@ -996,6 +1069,17 @@ def test_the_decoder_hands_over_good_records_before_it_raises(self): self.assertEqual(seen, [{"n": 1}]) self.assertEqual(caught.exception.code, "frame_too_large") + def test_the_decoder_hands_back_good_frames_before_it_refuses_a_bad_one(self): + """Same ordering guarantee, but the bad frame is junk rather than oversize.""" + decoder = wait_event.Decoder( + {"kind": "length_prefix", "prefix_bytes": 4, "byte_order": "big"}, 4096 + ) + records = decoder.feed(frame(b'{"n":1}') + frame(b"{not json")) + self.assertEqual(next(records), {"n": 1}) + with self.assertRaises(wait_event.SourceError) as caught: + next(records) + self.assertEqual(caught.exception.code, "invalid_json") + def test_a_line_decoder_also_hands_over_good_records_before_it_raises(self): decoder = wait_event.Decoder({"kind": "lines"}, 32) chunk = b'{"n": 1}\n' + b"x" * 64 + b"\n" @@ -1007,6 +1091,36 @@ def test_a_line_decoder_also_hands_over_good_records_before_it_raises(self): self.assertEqual(seen, [{"n": 1}]) self.assertEqual(caught.exception.code, "frame_too_large") + def test_an_explicit_empty_body_path_reads_the_whole_record_as_the_body(self): + """[] and an omitted key are the same request: the record is the body.""" + envelope = wait_event.parse_envelope( + {"match_fields": {}, "body_path": []}, "source.event_envelope" + ) + self.assertEqual(envelope["body_path"], []) + snapshot = wait_event.parse_snapshot( + { + "request": {"kind": "req"}, + "response_match_fields": {"kind": "res"}, + "body_path": [], + }, + "source.snapshot", + ) + self.assertEqual(snapshot["body_path"], []) + + def test_a_path_that_selects_a_value_still_may_not_be_empty(self): + """Accepting an empty body_path must not loosen paths that select a value.""" + for where in ("subject_path", "status_path"): + match = { + "subject_path": ["workflowId"], + "subject": "wf-1", + "status_path": ["status"], + "terminal_statuses": ["completed"], + } + match[where] = [] + with self.assertRaises(wait_event.SpecError) as caught: + wait_event.parse_match(match) + self.assertEqual(caught.exception.code, "spec_type") + if __name__ == "__main__": unittest.main(verbosity=2) diff --git a/tests/test_cat_mode.py b/tests/test_cat_mode.py index 5a71d86dc..8eca54c7b 100644 --- a/tests/test_cat_mode.py +++ b/tests/test_cat_mode.py @@ -703,6 +703,21 @@ def test_remote_install_does_not_let_install_sh_eat_the_script(self): source = handle.read() self.assertIn("./install.sh > /tmp/catstack-install.log 2>&1