diff --git a/dag/gunbc/spark/serving_observation_transaction.dag b/dag/gunbc/spark/serving_observation_transaction.dag index 2e9a888ad6f..39094d3208c 100644 --- a/dag/gunbc/spark/serving_observation_transaction.dag +++ b/dag/gunbc/spark/serving_observation_transaction.dag @@ -69,6 +69,8 @@ type SparkServingProbeRequest | ModelBlobDigestsRequest { digests: List } | UserUnitReadRequest | UserUnitIsEnabledRequest + | RunningDefinitionRequest + | InvocationIdentityRequest | RuntimeRootRequest | RuntimeExecutableRequest | RuntimeExecutableBitRequest @@ -95,6 +97,8 @@ fn spark_serving_probe_request_leg(request: SparkServingProbeRequest) -> SparkSe ModelBlobDigestsRequest { digests: _ } => ProbeModelBlobDigests UserUnitReadRequest => ProbeUserUnitRead UserUnitIsEnabledRequest => ProbeUserUnitIsEnabled + RunningDefinitionRequest => ProbeRunningDefinition + InvocationIdentityRequest => ProbeInvocationIdentity RuntimeRootRequest => ProbeRuntimeRoot RuntimeExecutableRequest => ProbeRuntimeExecutable RuntimeExecutableBitRequest => ProbeRuntimeExecutableBit @@ -122,6 +126,8 @@ type SparkServingProbeLeg | ProbeModelBlobDigests | ProbeUserUnitRead | ProbeUserUnitIsEnabled + | ProbeRunningDefinition + | ProbeInvocationIdentity | ProbeRuntimeRoot | ProbeRuntimeExecutable | ProbeRuntimeExecutableBit @@ -148,6 +154,8 @@ fn spark_serving_probe_leg_key(leg: SparkServingProbeLeg) -> NonEmptyStr { ProbeModelBlobDigests => "model_blob_digests" as NonEmptyStr ProbeUserUnitRead => "user_unit_read" as NonEmptyStr ProbeUserUnitIsEnabled => "user_unit_is_enabled" as NonEmptyStr + ProbeRunningDefinition => "running_definition" as NonEmptyStr + ProbeInvocationIdentity => "invocation_identity" as NonEmptyStr ProbeRuntimeRoot => "runtime_root" as NonEmptyStr ProbeRuntimeExecutable => "runtime_executable" as NonEmptyStr ProbeRuntimeExecutableBit => "runtime_executable_bit" as NonEmptyStr @@ -590,13 +598,20 @@ fn spark_serving_transaction_duplicated_legs( // only. An absent capture must not be left for whichever downstream reader happens to want it // first, because that reader will read Absent and honestly report unobserved -- and the whole host // will look measured-and-empty rather than never-acquired. +// +// FOLLOW-UP TRIGGER: the receipt witnesses still spell the evaluated executed-count snapshots as +// numeric literals. Before another probe leg joins this roster, expose the roster cardinality to +// those witnesses and assert the relation `executed == required - refused` instead. The present +// literals are honest for this closed population, but a third added leg can otherwise make the +// oracle stale until a person recounts it; deriving the relation from this authority removes that +// hand-maintained second count. fn spark_serving_required_probe_legs() -> List { [ ProbeLoadState, ProbeIsActive, ProbeVersion, ProbeTags, ProbeNvidiaDriver, ProbeRuntimeReceipt, ProbeActiveEnter, ProbeDefaultTargetEnter, ProbeUptime, ProbeGrants, ProbePasswd, ProbeAuthorizedKeys, ProbeLinger, ProbeModelManifest, ProbeModelBlobsListing, ProbeModelManifestDigest, ProbeModelBlobDigests, - ProbeUserUnitRead, ProbeUserUnitIsEnabled, + ProbeUserUnitRead, ProbeUserUnitIsEnabled, ProbeRunningDefinition, ProbeInvocationIdentity, ProbeRuntimeRoot, ProbeRuntimeExecutable, ProbeRuntimeExecutableBit, ProbeRuntimeRootListing, ] } diff --git a/dag/gunbc/spark/serving_observe.dag b/dag/gunbc/spark/serving_observe.dag index 08abf16a655..dd0aca670c4 100644 --- a/dag/gunbc/spark/serving_observe.dag +++ b/dag/gunbc/spark/serving_observe.dag @@ -1,6 +1,9 @@ module gunbc.spark.serving_observe import std.types { String, List, NonEmptyStr, Int } +import std.decl_ref { DeclarationRef, WholeDeclaration } +import std.disposition { Disposition, Scaffold, SingleAuthority } +import std.dissolution { DissolutionCondition, unbound_dissolution } import std.algebra { trim } import std.content_hash { serialize_content_hash, as_content_hash_cryptographic, sha256_hex_digest } import extdeps.languages.json.parse { @@ -33,6 +36,13 @@ import gunbc.spark.serving_terminal_health { spark_serving_served_proof, spark_serving_residency_from_capture, } +import gunbc.spark.serving_realization { + SparkServingRunningDefinition, + RunningDefinitionObserved, + RunningDefinitionUnstamped, + RunningDefinitionNotRunning, + RunningDefinitionUnobserved, +} import gunbc.ollama_runtime_bundle { ModelStoreObservation, ModelStoreAbsent, @@ -147,7 +157,7 @@ import gunbc.spark.serving_observation_transaction { ProbeRuntimeReceipt, ProbeActiveEnter, ProbeDefaultTargetEnter, ProbeUptime, ProbeGrants, ProbePasswd, ProbeAuthorizedKeys, ProbeLinger, ProbeModelManifest, ProbeModelBlobsListing, ProbeModelManifestDigest, ProbeModelBlobDigests, - ProbeUserUnitRead, ProbeUserUnitIsEnabled, + ProbeUserUnitRead, ProbeUserUnitIsEnabled, ProbeRunningDefinition, ProbeInvocationIdentity, ProbeRuntimeRoot, ProbeRuntimeExecutable, ProbeRuntimeExecutableBit, ProbeRuntimeRootListing, SparkServingProbeRequest, SparkServingProbeOutcome, @@ -159,7 +169,7 @@ import gunbc.spark.serving_observation_transaction { RuntimeReceiptRequest, ActiveEnterRequest, DefaultTargetEnterRequest, UptimeRequest, GrantsRequest, PasswdRequest, AuthorizedKeysRequest, LingerRequest, ModelManifestRequest, ModelBlobsListingRequest, ModelManifestDigestRequest, ModelBlobDigestsRequest, - UserUnitReadRequest, UserUnitIsEnabledRequest, + UserUnitReadRequest, UserUnitIsEnabledRequest, RunningDefinitionRequest, InvocationIdentityRequest, RuntimeRootRequest, RuntimeExecutableRequest, RuntimeExecutableBitRequest, RuntimeRootListingRequest, } import gunbc.spark.managed_access_bootstrap { SparkGrantOutcome, GrantEnumerated, GrantAbsent, GrantProbeRefused, spark_managed_executor_login, spark_managed_principal_home } @@ -249,6 +259,100 @@ fn spark_serving_active_probe_words() -> List { ] } +data spark_serving_running_definition_probe_shell_scaffold: Disposition = Scaffold { + dissolves_to: SingleAuthority, + bind: DeclarationRef { + module_path: "gunbc.spark.serving_observe", + decl_name: "spark_serving_running_definition_probe_argv", + field: WholeDeclaration, + } +} + +data spark_serving_running_definition_probe_shell_dissolution_trigger: DissolutionCondition = unbound_dissolution( + description: "dissolve-on: spark_serving_running_definition_probe_argv -- anti-race read/compare/read is temporarily a bash -lc medium-as-string realization. DISSOLVES WHEN orchestration can sequence typed argv observation legs, compare their typed captures inside one transaction, and refuse when the selected process generation moves, so this bracket is emitted without a hand-authored shell program." as NonEmptyStr, +) + +// READ THE DEFINITION IDENTITY FROM THE PROCESS SYSTEMD IS ACTUALLY SUPERVISING. +// +// Neither the unit file nor `systemctl show ExecStart` answers this question after daemon-reload: +// both describe the new definition while the old process can remain alive. MainPID selects the +// running invocation and /proc//environ retains the identity stamped into that invocation at +// exec. The PID AND Linux process start time are read around the environment so a concurrent +// restart refuses instead of joining bytes from one invocation to the manager state of another; +// the start-time bracket also catches the rarer case where the numeric PID is reused inside the +// window. +// +// A SUCCESSFUL read with no marker prints `absent`: that is positive evidence that this invocation +// cannot be established as having started from the current stamped definition. It says nothing +// chronological -- a foreign or manually started process can also be unstamped. It is not +// unobservable. Only inability to read procfs or a malformed/moving MainPID refuses. +fn spark_serving_running_definition_probe_argv() -> List { + [ + "bash", "-lc", + join([ + "pid_before=$(systemctl --user show ", spark_serving_unit_name, + " --property=MainPID --value) || exit $?; ", + "case \"$pid_before\" in ''|*[!0-9]*|0) echo 'running-definition: MainPID did not name a running process' >&2; exit 65;; esac; ", + "start_before=$(awk '\{sub(/^.*\\) /, \"\"); print $20\}' /proc/$pid_before/stat) || exit $?; ", + "identity=$(awk -v 'RS=\\0' -F= '$1 == \"GUNBC_SPARK_SERVING_DEFINITION_IDENTITY\" \{ print \"present:\" substr($0, index($0, \"=\") + 1); found=1 \} END \{ if (!found) print \"absent\" \}' /proc/$pid_before/environ) || exit $?; ", + "pid_after=$(systemctl --user show ", spark_serving_unit_name, + " --property=MainPID --value) || exit $?; ", + "start_after=$(awk '\{sub(/^.*\\) /, \"\"); print $20\}' /proc/$pid_after/stat) || exit $?; ", + "test \"$pid_before:$start_before\" = \"$pid_after:$start_after\" || \{ echo 'running-definition: MainPID or process start time moved during observation' >&2; exit 67; \}; ", + "printf '%s\\n' \"$identity\"", + ], ""), + ] +} + +data spark_serving_invocation_identity_probe_shell_scaffold: Disposition = Scaffold { + dissolves_to: SingleAuthority, + bind: DeclarationRef { + module_path: "gunbc.spark.serving_observe", + decl_name: "spark_serving_invocation_identity_probe_argv", + field: WholeDeclaration, + } +} + +data spark_serving_invocation_identity_probe_shell_dissolution_trigger: DissolutionCondition = unbound_dissolution( + description: "dissolve-on: spark_serving_invocation_identity_probe_argv -- supervisor-token and process-generation read/compare/read is temporarily a bash -lc medium-as-string realization. DISSOLVES WHEN orchestration can sequence typed argv observation legs, compare their typed captures inside one transaction, and refuse when identity, process generation, environment readability, or activity moves, so this bracket is emitted without a hand-authored shell program." as NonEmptyStr, +) + +// SYSTEMD'S INVOCATION ID IS THE GENERATION IDENTITY; PID AND START TIME ARE ITS BRACKET. +// +// InvocationID is minted by the supervisor for every activation, so two restarts from the same +// definition cannot compare equal. PID:start-time is still read on both sides: it binds the opaque +// supervisor token to the exact process whose environment is readable and whose unit is active, +// and prevents PID reuse or a concurrent restart from returning a plausible identity for a mixed +// observation. The environment is deliberately read here even though this carrier does not expose +// its definition stamp: otherwise the invocation identity would not identify the process about +// which the sibling running-definition probe can speak. +fn spark_serving_invocation_identity_probe_argv() -> List { + [ + "bash", "-lc", + join([ + "active_before=$(systemctl --user is-active ", spark_serving_unit_name, ") || exit $?; ", + "test \"$active_before\" = active || exit 65; ", + "pid_before=$(systemctl --user show ", spark_serving_unit_name, + " --property=MainPID --value) || exit $?; ", + "case \"$pid_before\" in ''|*[!0-9]*|0) echo 'invocation-identity: MainPID did not name a running process' >&2; exit 65;; esac; ", + "start_before=$(awk '\{sub(/^.*\\) /, \"\"); print $20\}' /proc/$pid_before/stat) || exit $?; ", + "invocation_before=$(systemctl --user show ", spark_serving_unit_name, + " --property=InvocationID --value) || exit $?; ", + "test -n \"$invocation_before\" || \{ echo 'invocation-identity: systemd returned an empty InvocationID' >&2; exit 66; \}; ", + "test -r /proc/$pid_before/environ || \{ echo 'invocation-identity: process environment is unreadable' >&2; exit 66; \}; ", + "awk -v 'RS=\\0' '\{ seen=1 \}' /proc/$pid_before/environ >/dev/null || exit $?; ", + "pid_after=$(systemctl --user show ", spark_serving_unit_name, + " --property=MainPID --value) || exit $?; ", + "start_after=$(awk '\{sub(/^.*\\) /, \"\"); print $20\}' /proc/$pid_after/stat) || exit $?; ", + "invocation_after=$(systemctl --user show ", spark_serving_unit_name, + " --property=InvocationID --value) || exit $?; ", + "active_after=$(systemctl --user is-active ", spark_serving_unit_name, ") || exit $?; ", + "test \"$pid_before:$start_before:$invocation_before:$active_before\" = \"$pid_after:$start_after:$invocation_after:$active_after\" || \{ echo 'invocation-identity: process generation or activity moved during observation' >&2; exit 67; \}; ", + "printf '%s\\n' \"$invocation_before\"", + ], ""), + ] +} + data spark_serving_version_probe_words: List = [ "curl", "-sf", "http://127.0.0.1:11434/api/version", ] @@ -849,6 +953,8 @@ fn spark_serving_acquire_observation_transaction( spark_serving_captured_leg(context: context, host: host, request: ModelManifestDigestRequest), spark_serving_captured_leg(context: context, host: host, request: UserUnitReadRequest), spark_serving_captured_leg(context: context, host: host, request: UserUnitIsEnabledRequest), + spark_serving_captured_leg(context: context, host: host, request: RunningDefinitionRequest), + spark_serving_captured_leg(context: context, host: host, request: InvocationIdentityRequest), spark_serving_captured_leg(context: context, host: host, request: RuntimeRootRequest), spark_serving_captured_leg(context: context, host: host, request: RuntimeExecutableRequest), spark_serving_captured_leg(context: context, host: host, request: RuntimeExecutableBitRequest), @@ -901,6 +1007,8 @@ fn spark_serving_probe_request_argv(request: SparkServingProbeRequest) -> List spark_serving_user_unit_read_argv() UserUnitIsEnabledRequest => spark_serving_user_unit_is_enabled_argv() + RunningDefinitionRequest => spark_serving_running_definition_probe_argv() + InvocationIdentityRequest => spark_serving_invocation_identity_probe_argv() RuntimeRootRequest => spark_serving_runtime_root_probe_argv() RuntimeExecutableRequest => spark_serving_runtime_executable_probe_argv() RuntimeExecutableBitRequest => spark_serving_runtime_executable_bit_probe_argv() @@ -969,6 +1077,8 @@ fn spark_serving_request_prerequisite_refusal( ModelManifestDigestRequest => none UserUnitReadRequest => none UserUnitIsEnabledRequest => none + RunningDefinitionRequest => none + InvocationIdentityRequest => none RuntimeRootRequest => none RuntimeExecutableRequest => none RuntimeExecutableBitRequest => none @@ -1952,6 +2062,141 @@ fn classify_is_active_result(result: SshSessionExecResult) -> SparkServingObserv } } +// WHICH DEFINITION THE RUNNING INVOCATION WAS STARTED FROM is a separate observation from unit +// activity. `is-active` establishes only running/not-running; this carrier joins that answer to the +// identity read from the process selected by MainPID. Keeping the subjects separate prevents an +// active process from becoming "current" merely because its file or the manager's loaded definition +// advanced after exec. +fn spark_serving_running_definition_from_probe_results( + is_active: SshSessionExecResult?, + running_definition: SshSessionExecResult?, +) -> SparkServingRunningDefinition { + match is_active { + Absent => RunningDefinitionUnobserved { + cause: "unit activity probe did not run, so no running invocation can be selected" as NonEmptyStr, + } + Present { value: active } => + if active.success && trim(s: active.stdout) == "active" { + match running_definition { + Absent => RunningDefinitionUnobserved { + cause: "unit is active but the running-definition probe did not run" as NonEmptyStr, + } + Present { value: definition } => + if definition.success && trim(s: definition.stdout) == "absent" { + RunningDefinitionUnstamped + } else if definition.success && starts_with(s: trim(s: definition.stdout), prefix: "present:") { + let wire = trim(s: definition.stdout) + let identity = substring(s: wire, start: 8, end: length(wire)) + if identity == "" { + RunningDefinitionUnobserved { + cause: "running process carried an empty definition identity" as NonEmptyStr, + } + } else { + RunningDefinitionObserved { definition_identity: identity as NonEmptyStr } + } + } else { + RunningDefinitionUnobserved { + cause: join([ + "running-definition probe refused exit=", to_string(definition.exit_code), + " stderr=", definition.stderr, + ], "") as NonEmptyStr, + } + } + } + } else if active.exit_code == 3 || active.exit_code == 4 || active.success { + RunningDefinitionNotRunning + } else { + RunningDefinitionUnobserved { + cause: join([ + "unit activity probe refused exit=", to_string(active.exit_code), + " stderr=", active.stderr, + ], "") as NonEmptyStr, + } + } + } +} + +// THIS PROJECTOR HAS NO PRODUCTION CONSUMER BY RULING, NOT BY OVERSIGHT. +// +// ProbeRunningDefinition is captured by the real observation transaction. Its only modeled consumer +// is gunbc.spark.serving_realization spark_serving_reactivation_remedy, a selector that nothing routes +// a production convergence through by ruling: DCH-0r-c must land the admission fence and drain first. +// Until that trigger lands the convergence path remains deliberately dark; this declaration does +// not make restart selectable from an unfenced observation. This observation discharges conjunct +// (i) of gunbc.recurring_failure_mode +// liveness_probe_read_as_currency's next-rung trigger; conjunct (ii), the guarded producer, remains +// DCH-0r-c. Being named by that trigger is a destination and NOT a production consumer or coverage. +fn spark_serving_running_definition_from_transaction( + txn: SparkServingObservationTransaction, +) -> SparkServingRunningDefinition { + spark_serving_running_definition_from_probe_results( + is_active: spark_serving_transaction_result(txn: txn, leg: ProbeIsActive), + running_definition: spark_serving_transaction_result(txn: txn, leg: ProbeRunningDefinition), + ) +} + +// INVOCATION GENERATION AND RUNNING DEFINITION HAVE DIFFERENT STABILITY. +// +// The definition identity is equal across restarts from one spec. This identity must differ across +// those restarts, so a post-reactivation observation can prove replacement rather than merely prove +// that both processes came from the same desired definition. Unobserved remains located; an empty +// or unreadable supervisor token must never be replaced with PID, definition identity, or a +// fabricated sentinel. +type SparkServingInvocationIdentity + = InvocationIdentityObserved { invocation_identity: NonEmptyStr } + | InvocationIdentityNotRunning + | InvocationIdentityUnobserved { cause: NonEmptyStr } + +fn spark_serving_invocation_identity_from_probe_results( + is_active: SshSessionExecResult?, + invocation: SshSessionExecResult?, +) -> SparkServingInvocationIdentity { + match is_active { + Absent => InvocationIdentityUnobserved { + cause: "unit activity probe did not run, so no serving invocation can be selected" as NonEmptyStr, + } + Present { value: active } => + if active.success && trim(s: active.stdout) == "active" { + match invocation { + Absent => InvocationIdentityUnobserved { + cause: "unit is active but the invocation-identity probe did not run" as NonEmptyStr, + } + Present { value: observed } => + if observed.success && trim(s: observed.stdout) != "" { + InvocationIdentityObserved { + invocation_identity: trim(s: observed.stdout) as NonEmptyStr, + } + } else { + InvocationIdentityUnobserved { + cause: join([ + "invocation-identity probe refused exit=", to_string(observed.exit_code), + " stderr=", observed.stderr, + ], "") as NonEmptyStr, + } + } + } + } else if active.exit_code == 3 || active.exit_code == 4 || active.success { + InvocationIdentityNotRunning + } else { + InvocationIdentityUnobserved { + cause: join([ + "unit activity probe refused exit=", to_string(active.exit_code), + " stderr=", active.stderr, + ], "") as NonEmptyStr, + } + } + } +} + +fn spark_serving_invocation_identity_from_transaction( + txn: SparkServingObservationTransaction, +) -> SparkServingInvocationIdentity { + spark_serving_invocation_identity_from_probe_results( + is_active: spark_serving_transaction_result(txn: txn, leg: ProbeIsActive), + invocation: spark_serving_transaction_result(txn: txn, leg: ProbeInvocationIdentity), + ) +} + // THE MODEL STORE PRODUCER -- the evidence behind the model-slot member. // // TWO LEGS, AND THE SECOND IS THE ONE THAT MATTERS. The manifest leg answers whether the reference diff --git a/dag/gunbc/spark/serving_unit_render.dag b/dag/gunbc/spark/serving_unit_render.dag index cdebc4a4637..39df3f7bfa9 100644 --- a/dag/gunbc/spark/serving_unit_render.dag +++ b/dag/gunbc/spark/serving_unit_render.dag @@ -1,6 +1,8 @@ module gunbc.spark.serving_unit_render import std.types { String, NonEmptyStr, FilePath } +import std.measure { token_count_value, positive_slot_count_value } +import std.dissolution { DissolutionCondition, unbound_dissolution } import std.content_hash { serialize_content_hash, content_hash_of_value } import std.bytes { bytes_octets, utf8_encode_bytes } import std.encoding { base64_encode, Standard } @@ -65,6 +67,79 @@ fn spark_serving_exec_start_command(spec: SparkServingUserUnitSpec) -> NonEmptyS ], "") as NonEmptyStr } +// THE DEFINITION IDENTITY CARRIED BY THE INVOCATION IT STARTS. +// +// systemd's loaded unit and the file on disk can both advance while an already-running process +// keeps the environment and command line from its earlier start. The process therefore carries +// the identity of the typed definition that created it. This is derived before rendering (and so +// is not a self-hash of a file containing its own hash), while still covering every typed unit-spec +// field. A later observer reads this value from /proc//environ; it never substitutes the +// file's current identity for the invocation's. +// THE CURRENT RUNG IS A COMPILE-TIME TRIPWIRE, NOT STRUCTURAL IMPOSSIBILITY. +// +// content_hash_of_value currently accepts only NonEmptyStr, so the typed spec must be projected to +// one canonical string wire. The total SparkServingUserUnitSpec reconstruction below makes adding a +// field fail compilation here and prompts an identity-domain decision; every serialized projection +// then reads from that reconstructed value. It remains mechanically possible for an author to add +// the field to the reconstruction but omit it from the join, so this converts a silent omission into +// a prompted one rather than making omission unwritable. String fields are base64-framed and field +// names fixed, so delimiters inside a field cannot make distinct covered specs collide. +// +// DISSOLVES WHEN the content-hash primitive accepts structured values and SparkServingUserUnitSpec +// can be hashed directly. That capability, not a record-fold artifact, makes every present and future +// field enter the domain by construction. This wire is intentionally NOT rendered unit-file content: +// those bytes contain the resulting identity and therefore cannot be their own hash domain. +fn spark_serving_user_unit_spec_identity_wire(spec: SparkServingUserUnitSpec) -> NonEmptyStr { + let identity_spec = SparkServingUserUnitSpec { + description: spec.description, + after: spec.after, + service_type: spec.service_type, + bind_host_port: spec.bind_host_port, + models_root_rel: spec.models_root_rel, + runtime_root_rel: spec.runtime_root_rel, + exec_subcommand: spec.exec_subcommand, + restart: spec.restart, + wanted_by: spec.wanted_by, + default_context: spec.default_context, + serving_slots: spec.serving_slots, + } + join([ + "description=", base64_encode(octets: bytes_octets(b: utf8_encode_bytes(s: identity_spec.description as String)), variant: Standard), "\n", + "after=", base64_encode(octets: bytes_octets(b: utf8_encode_bytes(s: identity_spec.after as String)), variant: Standard), "\n", + "service_type=", base64_encode(octets: bytes_octets(b: utf8_encode_bytes(s: identity_spec.service_type as String)), variant: Standard), "\n", + "bind_host_port=", base64_encode(octets: bytes_octets(b: utf8_encode_bytes(s: identity_spec.bind_host_port as String)), variant: Standard), "\n", + "models_root_rel=", base64_encode(octets: bytes_octets(b: utf8_encode_bytes(s: identity_spec.models_root_rel)), variant: Standard), "\n", + "runtime_root_rel=", base64_encode(octets: bytes_octets(b: utf8_encode_bytes(s: identity_spec.runtime_root_rel)), variant: Standard), "\n", + "exec_subcommand=", base64_encode(octets: bytes_octets(b: utf8_encode_bytes(s: identity_spec.exec_subcommand as String)), variant: Standard), "\n", + "restart=", base64_encode(octets: bytes_octets(b: utf8_encode_bytes(s: identity_spec.restart as String)), variant: Standard), "\n", + "wanted_by=", base64_encode(octets: bytes_octets(b: utf8_encode_bytes(s: identity_spec.wanted_by as String)), variant: Standard), "\n", + "default_context=", to_string(token_count_value(t: identity_spec.default_context)), "\n", + "serving_slots=", to_string(positive_slot_count_value(slots: identity_spec.serving_slots)), "\n", + ], "") as NonEmptyStr +} + +data spark_serving_user_unit_spec_identity_wire_dissolution_trigger: DissolutionCondition = unbound_dissolution( + description: "dissolve-on: spark_serving_user_unit_spec_identity_wire -- canonical string projection and total-record tripwire remain because content_hash_of_value accepts only NonEmptyStr. DISSOLVES WHEN the content-hash primitive accepts structured values and hashes SparkServingUserUnitSpec directly, placing every present and future field in the definition-identity domain by construction." as NonEmptyStr, +) + +fn spark_serving_user_unit_definition_identity(spec: SparkServingUserUnitSpec) -> NonEmptyStr { + serialize_content_hash( + hash: content_hash_of_value(value: spark_serving_user_unit_spec_identity_wire(spec: spec)), + ) as NonEmptyStr +} + +data spark_serving_definition_identity_environment_name: NonEmptyStr = "GUNBC_SPARK_SERVING_DEFINITION_IDENTITY" as NonEmptyStr + +fn spark_serving_definition_identity_environment_assignment( + spec: SparkServingUserUnitSpec, +) -> NonEmptyStr { + join([ + spark_serving_definition_identity_environment_name as String, + "=", + spark_serving_user_unit_definition_identity(spec: spec) as String, + ], "") as NonEmptyStr +} + // The runtime and model roots come from the install-paths authority rather than from the spec's // relative strings, so the unit and the installer cannot disagree about where the runtime lives. fn spark_serving_user_unit_file(spec: SparkServingUserUnitSpec) -> SystemdUnitFile { @@ -85,6 +160,7 @@ fn spark_serving_user_unit_file(spec: SparkServingUserUnitSpec) -> SystemdUnitFi Environment { assignment: ollama_num_parallel_env_assignment(slots: spec.serving_slots), }, + Environment { assignment: spark_serving_definition_identity_environment_assignment(spec: spec) }, WorkingDirectory { path: spark_serving_required_runtime_root() as NonEmptyStr }, ExecStart { command: spark_serving_exec_start_command(spec: spec) }, Restart { policy: spec.restart }, @@ -110,6 +186,16 @@ fn spark_serving_desired_user_unit_text() -> String { ) } +fn spark_serving_desired_user_unit_definition_identity() -> NonEmptyStr { + spark_serving_user_unit_definition_identity( + spec: spark_serving_user_unit_spec( + bind_host_port: spark_serving_desired_bind_host_port() as NonEmptyStr, + default_context: spark_serving_desired_default_context(), + serving_slots: spark_serving_desired_serving_slots(), + ), + ) +} + // ONE IDENTITY FOR UNIT CONTENT, used by the desired side and the observed side alike. // // The member's state IS the file's content, so its identity must be a function of bytes and nothing diff --git a/dag/test/claim/spark/spark_serving_hermetic_acquisition_fixture.dag b/dag/test/claim/spark/spark_serving_hermetic_acquisition_fixture.dag index 78817e1a52a..fb085249cff 100644 --- a/dag/test/claim/spark/spark_serving_hermetic_acquisition_fixture.dag +++ b/dag/test/claim/spark/spark_serving_hermetic_acquisition_fixture.dag @@ -74,6 +74,8 @@ import gunbc.spark.serving_observation_transaction { ModelBlobDigestsRequest, UserUnitReadRequest, UserUnitIsEnabledRequest, + RunningDefinitionRequest, + InvocationIdentityRequest, RuntimeRootRequest, RuntimeExecutableRequest, RuntimeExecutableBitRequest, @@ -194,7 +196,7 @@ fn hermetic_capture(request: SparkServingProbeRequest, out: String) -> SparkServ // THE COMPLETE POPULATION, BECAUSE A TRANSACTION MISSING ONE OF THEM IS NOT A TRANSACTION. // // Coverage is checked at identity grain against the required roster rather than by counting, so this -// enumerates the requests: a count would be satisfied by twenty-three copies of the wrong one. +// enumerates the requests: a count would be satisfied by twenty-five copies of the wrong one. fn hermetic_every_request() -> List admit_callers: [ decl_ref(module_path: "test.claim.spark.spark_serving_hermetic_acquisition_fixture", decl_name: "hermetic_complete_captures"), @@ -206,7 +208,7 @@ fn hermetic_every_request() -> List GrantsRequest, PasswdRequest, AuthorizedKeysRequest, LingerRequest, ModelManifestRequest, ModelBlobsListingRequest, ModelManifestDigestRequest, ModelBlobDigestsRequest { digests: ["deadbeef" as NonEmptyStr] }, - UserUnitReadRequest, UserUnitIsEnabledRequest, + UserUnitReadRequest, UserUnitIsEnabledRequest, RunningDefinitionRequest, InvocationIdentityRequest, RuntimeRootRequest, RuntimeExecutableRequest, RuntimeExecutableBitRequest, RuntimeRootListingRequest, ] diff --git a/dag/test/claim/spark/spark_serving_observation_transaction_witness_test.dag b/dag/test/claim/spark/spark_serving_observation_transaction_witness_test.dag index 7a975c7324a..3e0fcecc5ac 100644 --- a/dag/test/claim/spark/spark_serving_observation_transaction_witness_test.dag +++ b/dag/test/claim/spark/spark_serving_observation_transaction_witness_test.dag @@ -475,7 +475,7 @@ test fn the_receipt_line_says_which_provenance_arm_it_is() -> Bool { && !string_contains(s: hermetic, pattern: "observation-attempt=") } -// The transaction receipt reports coverage and the per-outcome shape, not a single total: 23 +// The transaction receipt reports coverage and the per-outcome shape, not a single total: 25 // captures is equally true of a transaction whose every leg the transport refused. test fn the_transaction_receipt_reports_coverage_and_outcome_shape() -> Bool { match txn_of(captures: complete_captures()) { @@ -486,7 +486,7 @@ test fn the_transaction_receipt_reports_coverage_and_outcome_shape() -> Bool { ) string_contains(s: line, pattern: "missing=0") && string_contains(s: line, pattern: "duplicated=0") - && string_contains(s: line, pattern: "executed=23") + && string_contains(s: line, pattern: "executed=25") && string_contains(s: line, pattern: "transport-refused=0") && string_contains(s: line, pattern: "prerequisite-refused=0") && string_contains(s: line, pattern: concat("observation-attempt=", concat(hermetic_fleet_attempt_raw, "/srv5/plan"))) @@ -536,7 +536,7 @@ test fn a_transport_refused_leg_moves_the_receipt_counts() -> Bool { let line = spark_serving_provenance_receipt_line( p: spark_serving_transaction_provenance(txn: t), ) - string_contains(s: line, pattern: "executed=22") + string_contains(s: line, pattern: "executed=24") && string_contains(s: line, pattern: "transport-refused=1") && string_contains(s: line, pattern: "missing=0") } @@ -581,7 +581,7 @@ test fn the_report_body_carries_the_acquisition_receipt_for_an_acquired_host() - Absent => false Present { value: body } => string_contains(s: body, pattern: "spark-serving-provenance host=srv5 acquired") - && string_contains(s: body, pattern: "executed=23") + && string_contains(s: body, pattern: "executed=25") && string_contains(s: body, pattern: "missing=0") && string_contains(s: body, pattern: "duplicated=0") } @@ -613,9 +613,9 @@ test fn changed_capture_outcomes_move_the_line_the_report_emits() -> Bool { match acquired_observation_body(captures: swap_is_active(with: refused)) { Absent => false Present { value: body } => - string_contains(s: body, pattern: "executed=22") + string_contains(s: body, pattern: "executed=24") && string_contains(s: body, pattern: "transport-refused=1") - && !string_contains(s: body, pattern: "executed=23") + && !string_contains(s: body, pattern: "executed=25") } } diff --git a/dag/test/claim/spark/spark_serving_observe_witness_test.dag b/dag/test/claim/spark/spark_serving_observe_witness_test.dag index edcc3709c7d..8df4fe2f57d 100644 --- a/dag/test/claim/spark/spark_serving_observe_witness_test.dag +++ b/dag/test/claim/spark/spark_serving_observe_witness_test.dag @@ -74,7 +74,20 @@ import gunbc.spark.serving_observe { spark_serving_boot_provenance_from_probes, SparkServingHostObservation, spark_serving_observation_report_text, + spark_serving_running_definition_from_probe_results, + spark_serving_invocation_identity_probe_argv, + spark_serving_invocation_identity_from_probe_results, + InvocationIdentityObserved, + InvocationIdentityNotRunning, + InvocationIdentityUnobserved, } +import gunbc.spark.serving_realization { + RunningDefinitionObserved, + RunningDefinitionUnstamped, + RunningDefinitionNotRunning, + RunningDefinitionUnobserved, +} + import gunbc.spark.serving_runtime_receipt { spark_serving_parse_runtime_receipt_json, spark_serving_runtime_receipt_field_from_probe, @@ -85,6 +98,140 @@ import test.claim.spark.spark_serving_hermetic_acquisition_fixture { hermetic_host_observation, } +fn running_definition_exec( + success: Bool, + exit_code: Int, + stdout: String, + stderr: String, +) -> SshSessionExecResult { + SshSessionExecResult { + success: success, exit_code: exit_code, stdout: stdout, stderr: stderr, + } +} + +test fn running_definition_identity_is_observed_from_the_invocation() -> Bool { + match spark_serving_running_definition_from_probe_results( + is_active: Present { value: running_definition_exec(success: true, exit_code: 0, stdout: "active\n", stderr: "") }, + running_definition: Present { value: running_definition_exec(success: true, exit_code: 0, stdout: "present:definition-a\n", stderr: "") }, + ) { + RunningDefinitionObserved { definition_identity: identity } => + identity == ("definition-a" as NonEmptyStr) + RunningDefinitionUnstamped => false + RunningDefinitionNotRunning => false + RunningDefinitionUnobserved { cause: _ } => false + } +} + +test fn invocation_identity_is_a_separate_observed_generation() -> Bool { + match spark_serving_invocation_identity_from_probe_results( + is_active: Present { value: running_definition_exec(success: true, exit_code: 0, stdout: "active\n", stderr: "") }, + invocation: Present { value: running_definition_exec(success: true, exit_code: 0, stdout: "9f6a-generation-a\n", stderr: "") }, + ) { + InvocationIdentityObserved { invocation_identity: identity } => + identity == ("9f6a-generation-a" as NonEmptyStr) + InvocationIdentityNotRunning => false + InvocationIdentityUnobserved { cause: _ } => false + } +} + +// The RED is authorable: returning a definition-derived constant for both observations makes this +// fail. Successive activations of one unchanged definition must remain distinguishable. +test fn successive_invocation_identities_can_differ_without_definition_drift() -> Bool { + let active = Present { value: running_definition_exec(success: true, exit_code: 0, stdout: "active\n", stderr: "") } + let first = spark_serving_invocation_identity_from_probe_results( + is_active: active, + invocation: Present { value: running_definition_exec(success: true, exit_code: 0, stdout: "generation-a\n", stderr: "") }, + ) + let second = spark_serving_invocation_identity_from_probe_results( + is_active: active, + invocation: Present { value: running_definition_exec(success: true, exit_code: 0, stdout: "generation-b\n", stderr: "") }, + ) + match first { + InvocationIdentityObserved { invocation_identity: a } => + match second { + InvocationIdentityObserved { invocation_identity: b } => a != b + InvocationIdentityNotRunning => false + InvocationIdentityUnobserved { cause: _ } => false + } + InvocationIdentityNotRunning => false + InvocationIdentityUnobserved { cause: _ } => false + } +} + +// The command itself is the executable control for all four generation properties. InvocationID is +// the value returned (not PID or definition identity); PID:start-time and InvocationID are bracketed +// across the environment read; and activity is checked on both sides of that same read. +test fn invocation_identity_probe_binds_supervisor_generation_process_environment_and_health() -> Bool { + let command = spark_serving_invocation_identity_probe_argv().skip(n: 2).first() + string_contains(s: command, pattern: "--property=InvocationID") + && string_contains(s: command, pattern: "/proc/$pid_before/environ") + && string_contains(s: command, pattern: "$pid_before:$start_before:$invocation_before:$active_before") + && string_contains(s: command, pattern: "$pid_after:$start_after:$invocation_after:$active_after") + && string_contains(s: command, pattern: "printf '%s\\n' \"$invocation_before\"") +} + +test fn invocation_identity_refuses_when_the_generation_cannot_be_observed() -> Bool { + match spark_serving_invocation_identity_from_probe_results( + is_active: Present { value: running_definition_exec(success: true, exit_code: 0, stdout: "active\n", stderr: "") }, + invocation: Present { value: running_definition_exec(success: false, exit_code: 67, stdout: "", stderr: "process generation moved") }, + ) { + InvocationIdentityUnobserved { cause: cause } => + string_contains(s: cause as String, pattern: "process generation moved") + InvocationIdentityObserved { invocation_identity: _ } => false + InvocationIdentityNotRunning => false + } +} + +test fn inactive_unit_has_no_invocation_identity() -> Bool { + match spark_serving_invocation_identity_from_probe_results( + is_active: Present { value: running_definition_exec(success: false, exit_code: 3, stdout: "inactive\n", stderr: "") }, + invocation: none, + ) { + InvocationIdentityNotRunning => true + InvocationIdentityObserved { invocation_identity: _ } => false + InvocationIdentityUnobserved { cause: _ } => false + } +} + +// The RED is authorable: changing the successful `absent` arm to Unobserved makes this fail, and +// would leave both currently-running legacy hosts indistinguishable from a permission failure. +test fn absent_identity_marker_is_positive_unstamped_evidence() -> Bool { + match spark_serving_running_definition_from_probe_results( + is_active: Present { value: running_definition_exec(success: true, exit_code: 0, stdout: "active\n", stderr: "") }, + running_definition: Present { value: running_definition_exec(success: true, exit_code: 0, stdout: "absent\n", stderr: "") }, + ) { + RunningDefinitionUnstamped => true + RunningDefinitionObserved { definition_identity: _ } => false + RunningDefinitionNotRunning => false + RunningDefinitionUnobserved { cause: _ } => false + } +} + +test fn unreadable_running_environment_is_located_unobservable() -> Bool { + match spark_serving_running_definition_from_probe_results( + is_active: Present { value: running_definition_exec(success: true, exit_code: 0, stdout: "active\n", stderr: "") }, + running_definition: Present { value: running_definition_exec(success: false, exit_code: 1, stdout: "", stderr: "Permission denied") }, + ) { + RunningDefinitionUnobserved { cause: cause } => + string_contains(s: cause as String, pattern: "Permission denied") + RunningDefinitionObserved { definition_identity: _ } => false + RunningDefinitionUnstamped => false + RunningDefinitionNotRunning => false + } +} + +test fn inactive_unit_has_no_running_definition() -> Bool { + match spark_serving_running_definition_from_probe_results( + is_active: Present { value: running_definition_exec(success: false, exit_code: 3, stdout: "inactive\n", stderr: "") }, + running_definition: none, + ) { + RunningDefinitionNotRunning => true + RunningDefinitionObserved { definition_identity: _ } => false + RunningDefinitionUnstamped => false + RunningDefinitionUnobserved { cause: _ } => false + } +} + data witness_required_release_asset_digest_hex: String = "79617139521db251c716d6229505d7530171ff4d68e476d0c22dabc15726a237" data witness_jetpack5_release_asset_digest_hex: String = "721af44a33380893cb35c2124832e64563e9abbd40b4c0280167110ff390dead" diff --git a/dag/test/claim/spark/spark_serving_unit_render_witness_test.dag b/dag/test/claim/spark/spark_serving_unit_render_witness_test.dag index 074060f59a6..88b6f0b34b5 100644 --- a/dag/test/claim/spark/spark_serving_unit_render_witness_test.dag +++ b/dag/test/claim/spark/spark_serving_unit_render_witness_test.dag @@ -14,6 +14,8 @@ import gunbc.spark.serving_unit_render { spark_serving_user_unit_content_identity, spark_serving_user_unit_content_wire, spark_serving_desired_user_unit_text, + spark_serving_user_unit_definition_identity, + spark_serving_definition_identity_environment_assignment, } import gunbc.spark.serving_desired { spark_serving_desired_bind_host_port, @@ -127,3 +129,21 @@ test fn the_desired_unit_is_rendered_at_the_declared_slot_count() -> Bool { == text_with_slots(slots: spark_serving_desired_serving_slots()) && spark_serving_desired_user_unit_text() != text_with_slots(slots: slots(n: 0)) } + +// The invocation marker is derived from the typed spec and rendered into the unit. Hashing the +// rendered bytes here would be circular because those bytes contain the marker itself; this control +// makes the well-founded direction executable and goes red if the marker stops reaching the file. +test fn the_unit_stamps_its_typed_definition_identity_into_the_invocation() -> Bool { + let spec = spec_with_window(window: token_count(count: 131072)) + let assignment = spark_serving_definition_identity_environment_assignment(spec: spec) + string_contains(s: spark_serving_user_unit_text(spec: spec), pattern: assignment as String) + && string_contains( + s: assignment as String, + pattern: spark_serving_user_unit_definition_identity(spec: spec) as String, + ) +} + +test fn the_running_definition_identity_changes_with_the_typed_spec() -> Bool { + spark_serving_user_unit_definition_identity(spec: spec_with_window(window: token_count(count: 8192))) + != spark_serving_user_unit_definition_identity(spec: spec_with_window(window: token_count(count: 131072))) +}