exception Defect = Loop.Defect exception Not_implemented = Loop.Not_implemented module M = Merge_queue type context = { env : Eio_unix.Stdenv.base; cwd : Fpath.t; directory : (Fpath.t, string) result; lookup : string -> string option; environ : string list; environment : (Repo.environment, string) result; caller : (Merge_queue.Caller.t, string) result; uid : int; say : Merge_queue.Progress.t -> unit; } let version = match Build_info.V1.version () with | Some v -> Build_info.V1.Version.to_string v | None -> "dev" let caller env ~lookup = let clock = Clock.v (Eio.Stdenv.clock env) (Eio.Stdenv.mono_clock env) in let user = match lookup "USER" with | Some u when u <> "" -> u | _ -> ( try Unix.getlogin () with Unix.Unix_error _ -> "unknown") in match M.Host.of_string (Unix.gethostname ()) with | Error (`Msg m) -> Error m | Ok host -> Ok { Merge_queue.Caller.user; host; terminal = Console_eio.is_tty Unix.stdin; tz_offset_s = Clock.tz_offset_s clock; } (* A variable of the environment, the last one when it is set twice, as getenv(3) reads it on the platforms mq runs on. *) let lookup_of environ = let vars = Array.fold_left (fun acc binding -> match String.index_opt binding '=' with | Some i -> ( String.sub binding 0 i, String.sub binding (i + 1) (String.length binding - i - 1) ) :: acc | None -> acc) [] environ in fun name -> List.assoc_opt name vars (* Why [dir] cannot be the directory a command runs in, when it cannot. *) let not_a_directory fs dir = match Eio.Path.kind ~follow:true Eio.Path.(fs / Fpath.to_string dir) with | `Directory -> None | `Not_found -> Fmt.kstr Option.some "%a does not exist" Fpath.pp dir | _ -> Fmt.kstr Option.some "%a is not a directory" Fpath.pp dir (* The current directory, then each -C in turn as git -C takes it: a relative one from the one before, each a directory to change to. *) let directory fs ~cwd dirs = List.fold_left (fun acc p -> Result.bind acc (fun dir -> let dir = Fpath.append dir p in match not_a_directory fs dir with | None -> Ok dir | Some why -> Error why)) (Ok cwd) dirs let context env ~environ ~uid ~cwd dirs = let fs = (Eio.Stdenv.fs env :> Eio.Fs.dir_ty Eio.Path.t) in let lookup = lookup_of environ in { env; cwd; directory = directory fs ~cwd dirs; lookup; environ = Array.to_list environ; environment = Repo.environment ~fs ~lookup ~uid; caller = caller env ~lookup; uid; say = (fun p -> Fmt.epr "%a@." M.Progress.pp p); } let todo what = raise (Not_implemented what) let run_name (caller : M.Caller.t) ~now = Fmt.str "%a:%d:%s" M.Host.pp caller.host (Unix.getpid ()) (Timestamp.Rfc3339.format now) (* {1 The queue a path names} *) type queue = { perform : Perform.t; location : Repo.location; rules : (M.Rules.t * M.Commit.t) option; unloadable : (M.Commit.t * string) option; registration : M.Registration.t option; process : [ `Generic | `Unix ] Eio.Process.mgr_ty Eio.Resource.t; } let refused ?branch ?next detail = M.Outcome.v ?branch ?next M.Outcome.Refused detail let unavailable ?branch detail = M.Outcome.v ?branch M.Outcome.Unavailable detail (* The repository's build slots: [mq.slots], else as many as the cores hold builds of the width a check runs at, which is every core while a check is given no width of its own. *) let slots ~fs ~mq_dir (local : M.Local_config.t) = let slots = match local.slots with | Some n -> n | None -> Build_slot.default_slots ~width:(Domain.recommended_domain_count ()) in Build_slot.v ~slots Eio.Path.(fs / Fpath.to_string (M.Slot.file ~mq_dir)) (* The shell's handles on one repository, once its local configuration is read. Nothing is written here: each store is made by its first write. *) (* Each check's median run over the last day of the journal, the window [mq stats] reads by default: the estimates the build slots are taken by (Order). A record the journal cannot read is left out. *) let estimates ~now journal = let day = Ptime.Span.of_int_s 86_400 in let since = Option.value ~default:Ptime.epoch (Ptime.sub_span now day) in let reading = Journal.since journal since in let flow = M.Flow.v { M.Flow.now; since; landings = []; moves = []; events = List.map (fun (r : Journal.record) -> r.event) reading.records; ranges = []; skipped = List.length reading.skipped; } in M.Order.estimates (List.filter_map (fun (c, (t : M.Flow.timing)) -> Option.map (fun m -> (c, m)) t.median_ms) flow.per_check) let perform_of ctx ~sw ~fs ~repo ~root ~caller local = let clock = Clock.v (Eio.Stdenv.clock ctx.env) (Eio.Stdenv.mono_clock ctx.env) in let mq_dir = Repo.mq_dir repo in let logs = Fpath.(mq_dir / "logs") in let pool = M.Local_config.pool local ~repository:root in let locks = match Run_lock.group [ Repo.common_dir repo ] with | Ok group -> group | Error reason -> failwith reason in let journal = Journal.v ~fs mq_dir ~source:("file://" ^ Fpath.to_string (Repo.common_dir repo)) in let checks = { Checks.fs; repo; pool; logs; logs_bound = Checks.default_bound; journal = (fun kind -> Journal.append journal (M.Event.v caller ~time:(Clock.now clock) kind)); clock; confine = Confine.v ~say:(fun line -> ctx.say (M.Progress.Line line)) ~process:(Eio.Stdenv.process_mgr ctx.env) ~fs ~lookup:ctx.lookup ~sleep:(Eio.Time.sleep (Eio.Stdenv.clock ctx.env)) ~repo ~pool local.M.Local_config.confine; path = Option.value ~default:"/usr/bin:/bin" (ctx.lookup "PATH"); home = Option.value ~default:"/" (ctx.lookup "HOME"); locks; slots = slots ~fs ~mq_dir local; owner = "mq"; estimates = (let now = Clock.now clock in Eio.Lazy.from_fun ~cancel:`Restart (fun () -> estimates ~now journal)); lease = Checks.lease (); } in { Perform.repo; journal; clock; caller; local; pool; checks; sw; fs; locks; run = None; repair = None; crontab = (fun () -> Scheduler.crontab ~fs ~path:(Option.value ~default:"" (ctx.lookup "PATH")) ~process: (Eio.Stdenv.process_mgr ctx.env :> [ `Generic | `Unix ] Eio.Process.mgr_ty Eio.Resource.t) ~sleep:(Eio.Time.sleep (Eio.Stdenv.clock ctx.env))); say = ctx.say; } (* Why no repository serves [path]. *) let discovery_failure path : Repo.discovery_error -> string = function | No_such_directory -> Fmt.str "%a does not exist" Fpath.pp path | Not_a_directory -> Fmt.str "%a is not a directory" Fpath.pp path | No_repository -> Fmt.str "%a is not in a Git repository" Fpath.pp path | No_checkout common -> Fmt.str "%a is in the bare repository %a, which has no main checkout" Fpath.pp path Fpath.pp common | Refused why -> why (* How a call failed to open its queue: refused, the path naming none or a key of the machine's holding a value the caller sets, with the command that fixes it; or unavailable, the queue unreadable. *) type failure = | Refuse of { why : string; next : string option } | Unavailable of string let refuse why = Refuse { why; next = None } let ( let* ) = Result.bind let config_of repo = Result.map_error (fun m -> Unavailable m) (Repo.configuration repo) let keys ~environment c = Result.map_error (fun (r : M.Local_config.refusal) -> Refuse { why = r.reason; next = r.next }) (M.Local_config.v ~env:(Repo.config_env environment) c) let network ctx = Repo.network ~net:(Eio.Stdenv.net ctx.env) ~mono:(Eio.Stdenv.mono_clock ctx.env) (* The repository a directory names, opened at its main checkout, with its keys (doc/commands.md, -C PATH). *) let open_repo ~sw ~network ~environment ~hooks (location : Repo.location) = let repo = Repo.v ~network ~sw ~environment ~hooks location.checkout in let* local = Result.bind (config_of repo) (keys ~environment) in Ok (repo, local) (* What mq schedule start recorded, and the keys a run reads under it: the rules branch the registration recorded, in place of mq.rules. *) let registered ~fs ~recorded repo local = match Scheduler.registered ~fs (Repo.mq_dir repo) with | Error m -> Error (Unavailable m) | Ok (Some r) when recorded -> Ok (Some r, M.Local_config.with_rules local r.M.Registration.rules) | Ok r -> Ok (r, local) (* Whether [cwd] is outside the repository's main checkout and every worktree of it: a command sent elsewhere by -C or GIT_DIR, whose caller is told which repository it opened (doc/commands.md, -C PATH). A directory outside the main checkout is inside a worktree when Git finds the repository's own common directory from it, as a walk from there finds it and not as GIT_DIR names it, so the worktrees are never listed. *) let outside environment cwd (location : Repo.location) = let real p = match Unix.realpath (Fpath.to_string p) with | r -> Fpath.to_dir_path (Fpath.v r) | exception Unix.Unix_error _ -> Fpath.to_dir_path p in let cwd' = real cwd in if Fpath.is_prefix (real location.checkout) cwd' then false else match Repo.discover (Repo.unsteered environment) cwd with | Ok here -> not (Fpath.equal (real here.common_dir) (real location.common_dir)) | Error _ -> true (* The rules at the tip of [local]'s rules branch. Rules that do not load there open the repository with no rules, and with the tip and why: each command answers them itself ({!initialised}, {!repairing}). A rules branch with no tip is unavailable. *) let rules_of repo (local : M.Local_config.t) = match Perform.read_rules repo local.rules with | Ok rules -> Ok (rules, None) | Error m -> ( match Repo.read_ref repo (M.Ref_name.branch local.rules) with | Some c -> Ok (None, Some (c, m)) | None -> Error (Unavailable m)) (* [repo] under the walls and the remote bound its rules set. *) let bounded repo = function | Some (r, _) -> let b = r.M.Rules.bounds in Repo.with_remote_bytes (Repo.with_wait_wall (Repo.with_hook_wall repo (M.Duration.of_seconds b.hook_wall)) (M.Duration.of_seconds b.wait_wall)) b.remote_bytes | None -> repo (* mq's own reference transactions start no hook mq init wrote: the hook's committed run is {!hook_committed}, in the work tree it would run in. *) let rec open_queue ?(recorded = false) ctx ~sw ~writes = let fs = (Eio.Stdenv.fs ctx.env :> Eio.Fs.dir_ty Eio.Path.t) in let process = (Eio.Stdenv.process_mgr ctx.env :> [ `Generic | `Unix ] Eio.Process.mgr_ty Eio.Resource.t) in let hooks = Repo.hooks ~committed:(fun work_tree lines -> hook_committed { ctx with directory = Ok work_tree } lines) ~process ~clock: (Eio.Stdenv.clock ctx.env :> float Eio.Time.clock_ty Eio.Resource.t) ~wall:(M.Duration.of_seconds M.Rules.default_bounds.hook_wall) () in let* caller = Result.map_error (fun m -> Unavailable m) ctx.caller in let* environment = Result.map_error (fun m -> Unavailable m) ctx.environment in let* dir = Result.map_error refuse ctx.directory in let* location = Result.map_error (fun e -> refuse (discovery_failure dir e)) (Repo.discover environment dir) in (* A repository whose configuration or rules mq cannot read is one it cannot serve: unavailable, in the words of what failed. A key of the machine's whose value is not one is the caller's to fix: refused, with the command that removes it (doc/configuration.md). *) let* repo, local, (rules, unloadable), registration = Loop.catching (fun () -> let* repo, local = open_repo ~sw ~network:(network ctx) ~environment ~hooks location in let* registration, local = registered ~fs ~recorded repo local in let* rules = rules_of repo local in Ok (repo, local, rules, registration)) ~on_failure:(fun m -> Error (Unavailable m)) in let announce = writes && outside environment ctx.cwd location in if announce then Fmt.kstr (fun s -> ctx.say (M.Progress.Line s)) "repository: %a" Fpath.pp location.checkout; let repo = bounded repo rules in let perform = perform_of ctx ~sw ~fs ~repo ~root:location.checkout ~caller local in Ok { perform; location; rules; unloadable; registration; process } (* Git aborts a transaction whose hook fails at preparing or prepared, and a move made outside mq is ordinary Git (doc/invariants.md, 1): the hook acts only once the transaction is committed, and whatever it meets, it writes nothing but the path the scheduler watches. *) and hook_committed ctx lines = match M.Registration.moved lines with | [] -> () | _ :: _ -> ( try Eio.Switch.run @@ fun sw -> match open_queue ctx ~sw ~writes:false with | Ok { rules = Some (rules, _); registration = Some r; perform; _ } when M.Registration.woken rules lines ~current:(fun name -> Option.map Git.Hash.to_hex (Repo.read_value perform.repo name)) -> Scheduler.wake ~fs:perform.fs r | Ok _ | Error _ -> () with Failure _ | Eio.Io _ | Unix.Unix_error _ | Sys_error _ -> ()) (* doc/runs.md: a target with branches queued and no run started within bounds.run_start is one event, and so is an edge left unserved for twice its interval. Nobody waits on a branch queued with --no-wait nor on an edge, so each command that writes journals the crossings its reading of the queue finds and the journal does not record yet. A queue it cannot read is left to the command's own reading, which names it. *) let journal_unserved q = match (q.rules, q.registration) with | None, _ | _, None -> Ok () | Some (rules, _), Some _ -> let p = q.perform in let run_start = M.Duration.of_seconds rules.M.Rules.bounds.run_start in let now = Perform.now p in let pins = match Result.bind (Repo.snapshot p.repo [ M.Ref_name.Prefix.queue ]) (Queue_read.collect (Queue_read.pin_of p.repo)) with | Ok pins -> pins | Error _ -> [] in let behind (e : M.Schedule.t) = e.behind > 0 in if M.Status.unserved ~run_start ~now pins [] = [] && not (List.exists behind (Queue_read.edges p rules [])) then Ok () else Loop.catching (fun () -> Ok (Journal.append_unrecorded p.journal ~segments:Queue_read.tail_segments (fun events -> List.map (fun (pin : M.Status.pin) -> Perform.event p ~branch:pin.entry.name (M.Event.Overdue { target = pin.entry.target; since = pin.queued })) (M.Status.unreported ~run_start ~now pins events) @ List.filter_map (fun (e : M.Schedule.t) -> Option.map (fun since -> Perform.event p ~branch:e.source (M.Event.Overdue { target = e.target; since })) e.served) (M.Schedule.unreported ~now (Queue_read.edges p rules events) events)))) ~on_failure:(fun m -> Error m) let with_queue ?branch ?recorded ?(unserved = true) ~writes ctx f = Eio.Switch.run @@ fun sw -> match open_queue ?recorded ctx ~sw ~writes with | Error (Refuse { why; next }) -> Error (refused ?branch ?next why) | Error (Unavailable m) -> Error (unavailable ?branch m) | Ok q when writes && unserved -> ( match journal_unserved q with | Ok () -> f q | Error m -> Error (unavailable ?branch m)) | Ok q -> f q let initialised ?branch q k = let rules_branch = q.perform.local.rules in match (q.rules, q.unloadable) with | Some rules, _ -> k rules | None, Some (_, why) -> Error (M.Outcome.diagnosed ?branch why) | None, None when M.Branch.equal rules_branch M.Branch.rules -> ( match Perform.past q.perform with | Never -> Error (refused ?branch ~next:"mq init" "mq init has not run in this repository") | Ran last -> Error (M.Init.gone ?branch ~rules_branch last)) | None, None -> Error (Fmt.kstr (refused ?branch) "mq.rules names %a, which is not a branch here" M.Branch.pp rules_branch) let repairing ?branch q ~target k = match (q.unloadable, target) with | Some (c, _), Some t when M.Branch.equal t q.perform.local.rules -> k { q with perform = { q.perform with repair = Some c } } (M.Rules.repairing, c) | _ -> initialised ?branch q (k q) let recorded_base q name = let p = q.perform in M.Description.recorded_base name ~upstream:(Perform.upstream p name) (Perform.description p name) let admitted ctx = with_queue ~writes:false ctx (fun q -> initialised q (fun _ -> Ok ()))