Skip to content
Closed
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
131 changes: 82 additions & 49 deletions dag/gunbc/fleet/fleet_converge_apply.dag
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,88 @@ import std.realization_reconcile {
NotConverged,
reconciliation_converged,
}
import gunbc.world_converge {
WorldObservation, WorldObserved, WorldUnobservable,
WorldDifference, WorldDiverged,
WorldActuation, WorldApplied, WorldActuationRefused,
WorldConvergenceHandler,
}

type HostConvergenceSubject { identity: HostIdentity }

type HostConvergenceContext {
policy: ConvergePolicy
transport: HostEffectTransport
}

type HostConvergenceDifference = HostConvergenceRequired

type HostConvergenceRefusal
= HostConvergencePolicyUnsound { reason: String }
| HostConvergenceUnknownIdentity { identity: HostIdentity }
| HostConvergenceApplyRefused { reason: String }

fn observe_host_convergence(
subject: HostConvergenceSubject,
context: HostConvergenceContext,
) -> WorldObservation<HostConverge, HostConvergenceRefusal> {
match context.policy {
ConvergePolicyUnsound { reason: reason } =>
WorldUnobservable { cause: HostConvergencePolicyUnsound { reason: reason } }
ConvergePolicyDerived { hosts: hosts } =>
match find_by_identity(hosts, identity: subject.identity, key: fn(host) { host.identity }) {
Present { value: host } => WorldObserved { value: host }
Absent => WorldUnobservable {
cause: HostConvergenceUnknownIdentity { identity: subject.identity },
}
}
}
}

fn host_convergence_differences(
_: HostConvergenceSubject,
_: HostConvergenceContext,
_: HostConverge,
) -> WorldDifference<HostConvergenceDifference> {
WorldDiverged { differences: [HostConvergenceRequired] }
}

fn apply_host_convergence(
_: HostConvergenceSubject,
context: HostConvergenceContext,
observed: HostConverge,
_: List<HostConvergenceDifference>,
) -> WorldActuation<Bool, HostConvergenceRefusal> {
match converge_apply_for_host(host_converge: observed, transport: context.transport) {
Absent => WorldActuationRefused {
cause: HostConvergenceUnknownIdentity { identity: observed.identity },
}
Present { value: result } =>
match result {
Converged { evidence: _, applied: _ } => WorldApplied { value: true }
NotConverged { reason: reason, applied: _ } => WorldActuationRefused {
cause: HostConvergenceApplyRefused { reason: reason },
}
}
}
}

fn unchanged_host_convergence_receipt(
_: HostConvergenceSubject,
_: HostConvergenceContext,
_: HostConverge,
) -> Bool {
false
}

type HostWorldConvergenceHandler = WorldConvergenceHandler<HostConvergenceSubject, HostConvergenceContext, HostConverge, HostConvergenceDifference, Bool, HostConvergenceRefusal>

data host_world_convergence_handler: HostWorldConvergenceHandler = WorldConvergenceHandler {
observe: observe_host_convergence,
differences: host_convergence_differences,
apply: apply_host_convergence,
unchanged_receipt: unchanged_host_convergence_receipt,
}

fn fleet_compute_host_for_identity(identity: HostIdentity) -> ComputeHost? {
find_by_identity(fleet_intent_known_hosts, identity: identity, key: fn(host) { host.identity })
Expand Down Expand Up @@ -58,55 +140,6 @@ fn converge_apply_for_host(
}
}

type FleetConvergeApplyResult =
FleetConvergeApplied { host_count: Int, all_converged: Bool }
| FleetConvergeApplyRefused { reason: String }

fn converge_apply(policy: ConvergePolicy, transport: HostEffectTransport) -> FleetConvergeApplyResult {
match policy {
ConvergePolicyUnsound { reason: reason } =>
FleetConvergeApplyRefused { reason: concat("fleet_converge_apply: policy unsound — ", reason) }
ConvergePolicyDerived { hosts: hosts } => {
let folded = fold(
hosts,
init: FleetConvergeApplied { host_count: 0, all_converged: true },
f: fn(acc, host_converge) {
match acc {
FleetConvergeApplyRefused { reason: _ } => acc
FleetConvergeApplied { host_count: n, all_converged: ok } =>
match converge_apply_for_host(host_converge: host_converge, transport: transport) {
Present { value: r } =>
FleetConvergeApplied {
host_count: n + 1,
all_converged: ok && reconciliation_converged(r: r),
}
Absent =>
FleetConvergeApplyRefused {
reason: concat(
"fleet_converge_apply: unknown fleet host identity ",
host_converge.identity as String
),
}
}
}
}
)
folded
}
}
}

fn converge_apply_fleet(transport: HostEffectTransport) -> FleetConvergeApplyResult {
converge_apply(policy: fleet_converge_policy(), transport: transport)
}

fn fleet_converge_thin_invocation(host: HostIdentity) -> ThinInvocation {
ThinInvocation { argv: ["gunbc", "converge", "--host", host as String] }
}

fn fleet_converge_apply_result_converged(r: FleetConvergeApplyResult) -> Bool {
match r {
FleetConvergeApplied { host_count: _, all_converged: ok } => ok
FleetConvergeApplyRefused { reason: _ } => false
}
}
114 changes: 96 additions & 18 deletions dag/gunbc/fleet/fleet_converge_plan_cli.dag
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,12 @@ import extdeps.filesystem.filesystem_io {
FilesystemEntryAbsent,
FilesystemEntryListed,
FilesystemEntryPresenceIndeterminate,
filesystem_file_observation,
filesystem_read_outcome,
FilesystemFileAbsent,
FilesystemFileRead,
FilesystemFileIndeterminate,
FilesystemFileObservationsDisagree,
}
import extdeps.exec.command { ArgvCommand, LocalExec }
import extdeps.tools.hostname { hostname_short_read_command }
Expand All @@ -23,6 +29,7 @@ import extdeps.posix.test_utility { test_regular_file_command }
import extdeps.languages.bash.invocation { bash_program_command }
import gunbc.command_runner { run_shell_command_capture, run_shell_commands, process_outcome_admitted }
import gunbc.host_axis_caps { CapMember, caps_managed_dropin_paths }
import gunbc.runner_unit { runner_unit_dropin_dir }
import gunbc.runner_slot_provision {
RunnerSlotMember,
RunnerSlotsObserved,
Expand Down Expand Up @@ -281,17 +288,71 @@ func observe_timer_members_wet() -> List<TimerRetirementMember>
)
}

func observe_cap_members_wet() -> List<CapMember>
type CapMembersObservation
= CapMembersObserved { members: List<CapMember> }
| CapMembersAbsent
| CapMembersUnobservable { path: NonEmptyStr, cause: String }

type CapMemberRead
= CapMemberReadObserved { member: CapMember }
| CapMemberReadAbsent
| CapMemberReadUnobservable { path: NonEmptyStr, cause: String }

fn cap_members_observation(reads: List<CapMemberRead>) -> CapMembersObservation {
fold(reads, init: CapMembersAbsent, f: fn(acc, read) {
match acc {
CapMembersUnobservable { path: _, cause: _ } => acc
CapMembersAbsent =>
match read {
CapMemberReadObserved { member: member } => CapMembersObserved { members: [member] }
CapMemberReadAbsent => CapMembersAbsent
CapMemberReadUnobservable { path: path, cause: cause } =>
CapMembersUnobservable { path: path, cause: cause }
}
CapMembersObserved { members: members } =>
match read {
CapMemberReadObserved { member: member } => CapMembersObserved {
members: append(members, items: [member]),
}
CapMemberReadAbsent => acc
CapMemberReadUnobservable { path: path, cause: cause } =>
CapMembersUnobservable { path: path, cause: cause }
}
}
})
}

func observe_cap_members_wet() -> CapMembersObservation
uses net: Network
{
flat_map(caps_managed_dropin_paths, path => {
let listed = Filesystem.List(path: runner_unit_dropin_dir as String)
let listing = filesystem_listing_observation(
directory: runner_unit_dropin_dir as String,
success: listed.success,
entries: listed.entries,
error: listed.error,
)
cap_members_observation(reads: map(caps_managed_dropin_paths, path => {
let read = Filesystem.Read(path: path as String)
if read.success {
[CapMember { dropin_path: path, content: read.content }]
} else {
[]
let name = match split(s: path as String, delimiter: "/").last() {
Present { value: value } => value
Absent => ""
}
})
match filesystem_file_observation(
listing: listing,
name: name,
path: path as String,
read: filesystem_read_outcome(content: read.content, success: read.success, error: read.error),
) {
FilesystemFileAbsent(_) => CapMemberReadAbsent
FilesystemFileRead { path: _, content: content } =>
CapMemberReadObserved { member: CapMember { dropin_path: path, content: content } }
FilesystemFileIndeterminate { cause: cause } =>
CapMemberReadUnobservable { path: path, cause: cause }
FilesystemFileObservationsDisagree { path: _, cause: cause } =>
CapMemberReadUnobservable { path: path, cause: cause }
}
}))
}

type FleetConvergeRequestObservation
Expand Down Expand Up @@ -346,17 +407,34 @@ func observe_fleet_converge_request_wet(scope_wire: String, host: NonEmptyStr) -
ApplySparkAxisAdmitted =>
match observe_runner_slot_members_wet() {
RunnerSlotsUnobserved { reason: why } => FleetConvergeRequestUnobserved { reason: why as String }
RunnerSlotsObserved { members: observed_slots } => FleetConvergeRequestObserved {
request: FullHostConverge {
host: host,
observed_timers: observe_timer_members_wet(),
observed_caps: observe_cap_members_wet(),
observed_spark: [],
observed_slots: observed_slots,
observed_fabric_cells: admit_fabric_cell_probe(probe: fabric_cell_probe_wet(host: host as HostIdentity)),
observed_activation_readiness: runner_activation_readiness_unobserved_no_transaction,
},
}
RunnerSlotsObserved { members: observed_slots } =>
match observe_cap_members_wet() {
CapMembersUnobservable { path: path, cause: cause } => FleetConvergeRequestUnobserved {
reason: join(["CapMemberUnobservable path=", path as String, " — ", cause], ""),
}
CapMembersAbsent => FleetConvergeRequestObserved {
request: FullHostConverge {
host: host,
observed_timers: observe_timer_members_wet(),
observed_caps: [],
observed_spark: [],
observed_slots: observed_slots,
observed_fabric_cells: admit_fabric_cell_probe(probe: fabric_cell_probe_wet(host: host as HostIdentity)),
observed_activation_readiness: runner_activation_readiness_unobserved_no_transaction,
},
}
CapMembersObserved { members: observed_caps } => FleetConvergeRequestObserved {
request: FullHostConverge {
host: host,
observed_timers: observe_timer_members_wet(),
observed_caps: observed_caps,
observed_spark: [],
observed_slots: observed_slots,
observed_fabric_cells: admit_fabric_cell_probe(probe: fabric_cell_probe_wet(host: host as HostIdentity)),
observed_activation_readiness: runner_activation_readiness_unobserved_no_transaction,
},
}
}
}
}
}
Expand Down
Loading
Loading