Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
54 changes: 54 additions & 0 deletions uraniborg/docs/automate_observation.md
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,7 @@ Optional fields are left out rather than set to `null`.
| :--- | :--- | :--- |
| `run_started` | `argv` (arguments, without the script name), `pid` | First event of every run |
| `step` | `step`, `state` (`started` / `finished` / `failed`), `device`\*, `duration_ms`\*\*, `message`\* | Around each phase (see below) |
| `step_progress` | `step`, `device`\*, `done`, `total` | Between a step's `started` and its `finished` / `failed`, for long steps (see [Step progress](#step-progress)) |
| `devices` | `devices` (list of `{serial, unauthorized, model?, product?, device?}`), `selected` (serials to observe), `missing` (requested via `--serial` but not connected) | Once, after listing devices |
| `device_started` | `device` | Before processing each selected device |
| `prompt` | `device`, `kind`, `message`, `expects_input` | The script is waiting for a person (see [Manual intervention](#manual-intervention)) |
Expand All @@ -123,6 +124,35 @@ pre-fetching is attempted) and `inclusion_proof_check`. A failed
`inclusion_proof_prefetch` is not fatal: verification falls back to fetching
entries on demand.

#### Step progress

Some steps take minutes. While they run, `step_progress` events report how far
they have got, e.g. to show "120 of 314 checked". Currently only
`inclusion_proof_check` reports progress: `done` is the number of APK splits
verified so far (found in the log or not) and `total` the number of splits to
verify (splits skipped for missing fields are not counted).

- Each event carries the same `step` and `device` as the step it belongs to,
and appears only between that step's `started` and `finished` / `failed`.
- The first event has `done` = `0`, once verification starts. For
`inclusion_proof_check`, that is after the package list has been read; if it
cannot be read, the step fails without any `step_progress`.
- `done` only increases, and `total` stays the same within a step.
- Events are throttled: at most one every 2 seconds, plus the first one and the
one where `done` reaches `total`, which is always sent. A step with nothing
to verify sends a single event with `done` = `total` = `0`.
- A step that is stopped early (Ctrl-C, `SIGTERM`, an error) ends with
`failed` without reaching `done` = `total`.

```json
{"v": 1, "ts": "2026-01-02T03:05:00.000Z", "type": "step", "step": "inclusion_proof_check", "state": "started", "device": "ABCDEF012345"}
{"v": 1, "ts": "2026-01-02T03:05:00.010Z", "type": "step_progress", "step": "inclusion_proof_check", "device": "ABCDEF012345", "done": 0, "total": 314}
{"v": 1, "ts": "2026-01-02T03:05:02.020Z", "type": "step_progress", "step": "inclusion_proof_check", "device": "ABCDEF012345", "done": 11, "total": 314}
...
{"v": 1, "ts": "2026-01-02T03:05:55.400Z", "type": "step_progress", "step": "inclusion_proof_check", "device": "ABCDEF012345", "done": 314, "total": 314}
{"v": 1, "ts": "2026-01-02T03:05:55.410Z", "type": "step", "step": "inclusion_proof_check", "state": "finished", "device": "ABCDEF012345", "duration_ms": 55410}
```

#### Manual intervention

Two situations need a person. Each is reported as a `prompt` event, and is
Expand Down Expand Up @@ -294,6 +324,30 @@ flags:
`preinstalled_packages.txt` instead of `packages.txt`, writing results to
`preinstalled_packages_with_inclusion_proof_signal.txt`.

### Batch Verification
All APK splits are verified by a single `verifier` run in batch mode
(`--payloads_path`, see the [verifier README](../../verifier_tools/verify/README.md)).
It fetches each log's checkpoint, and searches the log, once for all splits.
This keeps a cold cache (e.g. with `--no_prefetch`, or after pre-fetching
failed) from being filled by many verifiers downloading the same files at
once. Results are reported as the verifier finds them.

With a `verifier` built before batch mode existed, or for any splits a batch
run gives no result for (e.g. if it crashes), each split is verified by a
`verifier` run of its own instead, one at a time. That is much slower, so a
warning suggests rebuilding the `verifier`. The results file is the same
either way, including the order of packages and splits.

Log lines about individual splits (shown with `-D`) start with the package and
split they are about, e.g. `com.android.chrome [config.en]: ...`.

Ctrl-C stops the running verifier and starts no new ones. `SIGTERM` does the
same when `--events` is given (and always for `inclusion_proof_check.py`).
Without `--events`, `automate_observation.py` keeps the default `SIGTERM`
behaviour: the script exits at once, and a verifier that was already running
stops on its own: a per-split verifier after its split (about a second or
two), a batch verifier when it next tries to write a result.

### Pulling Pre-installed APKs Only
By default, `--pull-all-apks` downloads all packages listed in `packages.txt`.
You can pass `--pull-preinstalled-apks-only` to download only pre-installed
Expand Down
109 changes: 55 additions & 54 deletions uraniborg/scripts/python/automate_observation.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@

import inclusion_proof_check
import syscall_wrapper
import termination

AdbWrapper = syscall_wrapper.AdbWrapper
SyscallWrapper = syscall_wrapper.SyscallWrapper
Expand Down Expand Up @@ -227,6 +228,10 @@ def set_up_logging(args: argparse.Namespace) -> logging.Logger:

EVENTS_SCHEMA_VERSION = 1

# step_progress events are written at most this often (plus the first and
# the final one of each step).
STEP_PROGRESS_INTERVAL_SECONDS = 2.0

# Per-device outcomes. These are the same four outcomes the final log summary
# distinguishes.
STATUS_SUCCESS = "success"
Expand Down Expand Up @@ -273,30 +278,11 @@ def set_up_logging(args: argparse.Namespace) -> logging.Logger:
PROMPT_OUTCOME_STDIN_CLOSED = REASON_STDIN_CLOSED


class Terminated(BaseException):
"""Raised by the SIGTERM handler that main() installs while --events is on.

Derives from BaseException, like KeyboardInterrupt, so that the generic
`except Exception` handlers do not swallow it. SyscallWrapper only catches
KeyboardInterrupt, so this also propagates out of adb calls.
"""

def __init__(self):
super().__init__("received SIGTERM")


def _raise_terminated(signum, frame): # pylint: disable=unused-argument
# Ignore repeated SIGTERMs while cleaning up so that run_finished is still
# written; _die_by_sigterm() restores the default action afterwards.
signal.signal(signal.SIGTERM, signal.SIG_IGN)
raise Terminated()


def _interruption_error(e: BaseException) -> dict:
"""Describes an exception that is not an Exception, for device errors."""
if isinstance(e, KeyboardInterrupt):
return _error(REASON_INTERRUPTED, "Interrupted.")
if isinstance(e, Terminated):
if isinstance(e, termination.Terminated):
return _error(REASON_TERMINATED, "Terminated.")
return _error(REASON_UNEXPECTED_ERROR,
"{}: {}".format(type(e).__name__, e))
Expand Down Expand Up @@ -359,11 +345,13 @@ class EventEmitter:

def __init__(self, stream: Optional[TextIO] = None,
logger: Optional[logging.Logger] = None,
clock=None):
clock=None, monotonic=None):
self._stream = stream
self._logger = logger
self._clock = clock or (
lambda: datetime.datetime.now(datetime.timezone.utc))
# Used only to throttle step_progress events.
self._monotonic = monotonic or time.monotonic
self._run_finished = False
# serial -> {status, results_dir?}, in device_finished order.
self._finished_devices = {}
Expand Down Expand Up @@ -402,6 +390,31 @@ def device_finished(self, device: str, status: str,
self.emit("device_finished", device=device, status=status,
results_dir=results_dir, error=error)

def progress_reporter(
self, step: str, device: Optional[str] = None,
interval: float = STEP_PROGRESS_INTERVAL_SECONDS):
"""Returns a progress(done, total) callback emitting step_progress.

Meant to be called often (e.g. once per item); events are throttled:
the first call is always reported, and so is the one where done reaches
total, but in between at most one event is written per `interval`
seconds. A call repeating the last reported `done` is never reported.
"""
last = {"done": None, "at": None}

def report(done: int, total: int) -> None:
if self._stream is None or done == last["done"]:
return
now = self._monotonic()
if (last["at"] is not None and done != total and
now - last["at"] < interval):
return
last["done"], last["at"] = done, now
self.emit("step_progress", step=step, device=device, done=done,
total=total)

return report

@contextlib.contextmanager
def step(self, name: str, device: Optional[str] = None):
"""Brackets a phase with step started/finished/failed events.
Expand Down Expand Up @@ -1559,7 +1572,8 @@ def _observe_device(target_device, args: argparse.Namespace,
concurrency=args.cache_prefetch_concurrency,
timeout=args.cache_prefetch_timeout,
prefetch=False,
preinstalled_only=args.check_preinstalled_only):
preinstalled_only=args.check_preinstalled_only,
progress=events.progress_reporter("inclusion_proof_check", serial)):
progress.check_incomplete = True
# False means the check could not complete (bad input or unwritable
# output), not that some splits are absent from the log; per-split
Expand Down Expand Up @@ -1748,45 +1762,32 @@ def main():
# being written, turn it into an exception so the run can report itself,
# then die by SIGTERM anyway so the parent sees the usual exit status.
# Without --events, SIGTERM handling is left untouched.
previous_sigterm_handler = None
if events.enabled:
previous_sigterm_handler = signal.signal(signal.SIGTERM, _raise_terminated)

terminated = False
try:
exit_code = run(args, logger, events)
except Terminated:
events.finish_run(128 + signal.SIGTERM, error=_error(
REASON_TERMINATED, "Terminated by SIGTERM."))
terminated = True
except KeyboardInterrupt:
events.finish_run(130, error=_error(REASON_INTERRUPTED,
"Interrupted by user."))
raise
except BaseException as e:
events.finish_run(1, error=_error(
REASON_UNEXPECTED_ERROR, "{}: {}".format(type(e).__name__, e)))
raise
finally:
events.close()
if previous_sigterm_handler is not None:
signal.signal(signal.SIGTERM, previous_sigterm_handler)
with (termination.sigterm_raises() if events.enabled
else contextlib.nullcontext()):
try:
exit_code = run(args, logger, events)
except termination.Terminated:
events.finish_run(128 + signal.SIGTERM, error=_error(
REASON_TERMINATED, "Terminated by SIGTERM."))
terminated = True
except KeyboardInterrupt:
events.finish_run(130, error=_error(REASON_INTERRUPTED,
"Interrupted by user."))
raise
except BaseException as e:
events.finish_run(1, error=_error(
REASON_UNEXPECTED_ERROR, "{}: {}".format(type(e).__name__, e)))
raise
finally:
events.close()

if terminated:
_die_by_sigterm(logger)
termination.die_by_sigterm(logger)
return
if exit_code:
sys.exit(exit_code)


def _die_by_sigterm(logger):
"""Terminates this process with SIGTERM's default action."""
logger.error("Terminated by SIGTERM.")
signal.signal(signal.SIGTERM, signal.SIG_DFL)
os.kill(os.getpid(), signal.SIGTERM)
# Not reached unless SIGTERM is blocked; fall back to the conventional code.
sys.exit(128 + signal.SIGTERM)


if __name__ == "__main__":
main()
Loading
Loading