An expression that makes a function value keeps its module mapped, a process a merged program starts does not take the session's socket, and an evaluation's value is its own expression's and not one a restart resumed

This commit is contained in:
Joseph Ferano 2026-09-25 16:37:43 +07:00
parent 020392b064
commit 96dbcb1f2f
6 changed files with 216 additions and 61 deletions

View File

@ -283,6 +283,19 @@ let stop_reply t =
let stop_gen t : int option = Option.map fst (stop_reply t)
(* How many evaluated expressions the agent has queued, and the highest one
that has returned a value; [None] from an agent without the verb. *)
let calls t : (int * int) option =
match request t "calls" with
| exception Unix.Unix_error _ -> None
| text ->
(match String.split_on_char ' ' (String.trim text) with
| [ q; v ] ->
(match int_of_string_opt q, int_of_string_opt v with
| Some q, Some v -> Some (q, v)
| _ -> None)
| _ -> None)
(* The same stop with whose code it stopped in: [Some true] when the thread
was running an evaluated thunk, [Some false] when it was in the program's
own code, [None] from an agent that does not say. *)
@ -1253,6 +1266,10 @@ let eval_expr t ~code ~origin ~pause =
match Session.eval_expr ~origin ~pause t.session code with
| c ->
let before = match result t with Some (g, _) -> g | None -> 0L in
(* This expression's number among those the agent has queued: an
earlier one resumed by a restart can publish after this one is sent,
and the result counter alone would take its value for this one's. *)
let mine = Option.map (fun (q, _) -> q + 1) (calls t) in
(* Read here, beside [before], and for the same kind of reason: all
three are the "how things stood" half of a difference the wait below
measures. A program already sitting in a break when the request
@ -1406,9 +1423,16 @@ let eval_expr t ~code ~origin ~pause =
the sleep has to stay a sleep. *)
drain t;
let value () =
match result t with
| Some (g, v) when Int64.compare g before > 0 -> Some v
| _ -> None
let returned =
match mine, calls t with
| Some m, Some (_, v) -> v >= m
| _ -> true
in
if not returned then None
else
match result t with
| Some (g, v) when Int64.compare g before > 0 -> Some v
| _ -> None
in
match value () with
| Some v -> `Value v

View File

@ -5770,6 +5770,43 @@ let program ?(checks = true) ?(dev = false) ?(debug = false) ?(pnames = [])
String literals still have to come along: they are this module's own
constants, and omitting them is an undefined [@.str.N] at link time. *)
(* Whether an expression thunk makes a function value anywhere in its body or
in the clauses lifted out of it. Such a value's code address is in this
module — a lambda's body, or the thick wrapper a named function is handed
out through — and it may be stored anywhere, so the module must stay
mapped. Both backends ask this before marking a thunk's module
unloadable. *)
let thunk_makes_fn_values (p : Tast.program) name =
let mine = Hashtbl.create 8 in
Hashtbl.replace mine name ();
(* Lifted clauses nest: a lambda inside a lambda is lifted out of the
outer one's body, so the set grows until nothing new joins it. *)
let rec close () =
let grew = ref false in
List.iter
(fun (f : Tast.fn) ->
match f.Tast.fparent with
| Some q when Hashtbl.mem mine q && not (Hashtbl.mem mine f.Tast.name) ->
Hashtbl.replace mine f.Tast.name (); grew := true
| _ -> ())
p.Tast.fns;
if !grew then close ()
in
close ();
let found = ref false in
List.iter
(fun (f : Tast.fn) ->
if Hashtbl.mem mine f.Tast.name then
List.iter
(Tast.walk (fun (e : Tast.expr) ->
match e.Tast.e with
| Tast.FnAddr _ | Tast.Closure _ | Tast.Thicken _ -> found := true
| _ -> ()))
f.Tast.body)
p.Tast.fns;
!found
let redefinition ?(checks = true) ?(dev = false) ?(debug = false)
?(known = fun _ -> true) ?(retains = true)
?call ?(consts = []) ?(annotate = false) (p : Tast.program) ~fns
@ -6045,7 +6082,8 @@ let redefinition ?(checks = true) ?(dev = false) ?(debug = false)
into the result buffer, so nothing outside the module holds an address
inside it once the call has returned. Without this, clicking through
the frames of a break loop costs a permanent mapping per click. *)
if fns = [ fn ] && consts = [] && ((not retains) || m.nstr = 0) then
if fns = [ fn ] && consts = [] && ((not retains) || m.nstr = 0)
&& not (thunk_makes_fn_values p fn) then
Buffer.add_string m.out "\n@flan_reload_transient = global i8 1\n"
| None -> ()
end;

View File

@ -5621,7 +5621,7 @@ let redefinition ~checks ?(dev = true) ?(known = fun _ -> true)
function's registry names are not counted: the registry copies them. *)
(match call with
| Some fn
when fns = [ fn ] && consts = []
when fns = [ fn ] && consts = [] && not (Emit.thunk_makes_fn_values p fn)
&& ((not retains) || md.Emit.nstr = 0) ->
Buffer.add_string out
"\n\t.data\n\t.globl\tflan_reload_transient\n\

View File

@ -65,8 +65,9 @@ let send path line =
binary is the parent of every program it starts, which is the
--two-process shape. *)
let daemon_env path =
[| "FLAN_AGENT_SOCKET=" ^ path;
"FLAN_AGENT_OWNER=" ^ string_of_int (Unix.getpid ()) |]
let me = string_of_int (Unix.getpid ()) in
[| "FLAN_AGENT_SOCKET=" ^ path; "FLAN_AGENT_OWNER=" ^ me;
"FLAN_DEV_PARENT=" ^ me |]
let () =
match Sys.command "command -v clang > /dev/null 2>&1 && command -v llc > /dev/null 2>&1" with
@ -342,59 +343,71 @@ let () =
is not this program's, and it picks and announces a path of its own as
if nothing were set. The file standing in for the session's socket has
to still be the same file afterwards. *)
let stolen = tmp "stolen.sock" and serr = tmp "stolen.err" in
Out_channel.with_open_bin stolen (fun oc ->
output_string oc "the session's");
let senv =
Array.append aenv
[| "FLAN_AGENT_SOCKET=" ^ stolen; "FLAN_AGENT_OWNER=1" |]
let inherited ~shape ~owner =
let stolen = tmp "stolen.sock" and serr = tmp "stolen.err" in
Out_channel.with_open_bin stolen (fun oc ->
output_string oc "the session's");
let senv =
Array.append aenv
[| "FLAN_AGENT_SOCKET=" ^ stolen; "FLAN_AGENT_OWNER=" ^ owner |]
in
let s1 = ofd (tmp "stolen.out") and s2 = ofd serr in
let spid = Unix.create_process_env aexe [| aexe |] senv Unix.stdin s1 s2 in
Unix.close s1;
Unix.close s2;
let prefix = "flan agent: listening on " in
let sannounced () =
let text = In_channel.with_open_bin serr In_channel.input_all in
List.find_map
(fun l ->
if String.length l > String.length prefix
&& String.sub l 0 (String.length prefix) = prefix
then Some (String.sub l (String.length prefix)
(String.length l - String.length prefix))
else None)
(String.split_on_char '\n' text)
in
(match
if await (fun () -> sannounced () <> None) then sannounced () else None
with
| None ->
fail "%s: a program with someone else's FLAN_AGENT_SOCKET announced no \
socket of its own" shape;
(try Unix.kill spid Sys.sigkill with Unix.Unix_error _ -> ())
| Some p ->
if p = stolen then fail "%s: the inherited path was bound: %S" shape p;
if not (await (fun () -> Sys.file_exists p)) then
fail "%s: nothing was bound at the announced %S" shape p
else ignore (send p aso);
let reaped =
await ~ms:5000 (fun () ->
match Unix.waitpid [ Unix.WNOHANG ] spid with
| 0, _ -> false
| _ -> true)
in
if not reaped then begin
(try Unix.kill spid Sys.sigkill with Unix.Unix_error _ -> ());
fail "%s: the program with an inherited variable never finished" shape
end);
(match In_channel.with_open_bin stolen In_channel.input_all with
| "the session's" -> ()
| _ -> fail "%s: the inherited FLAN_AGENT_SOCKET's file was replaced" shape
| exception Sys_error _ ->
fail "%s: the inherited FLAN_AGENT_SOCKET's file was removed" shape);
List.iter (fun f -> try Sys.remove f with Sys_error _ -> ())
[ stolen; serr; tmp "stolen.out" ]
in
let s1 = ofd (tmp "stolen.out") and s2 = ofd serr in
let spid = Unix.create_process_env aexe [| aexe |] senv Unix.stdin s1 s2 in
Unix.close s1;
Unix.close s2;
let prefix = "flan agent: listening on " in
let sannounced () =
let text = In_channel.with_open_bin serr In_channel.input_all in
List.find_map
(fun l ->
if String.length l > String.length prefix
&& String.sub l 0 (String.length prefix) = prefix
then Some (String.sub l (String.length prefix)
(String.length l - String.length prefix))
else None)
(String.split_on_char '\n' text)
in
(match
if await (fun () -> sannounced () <> None) then sannounced () else None
with
| None ->
fail "a program with someone else's FLAN_AGENT_SOCKET announced no \
socket of its own";
(try Unix.kill spid Sys.sigkill with Unix.Unix_error _ -> ())
| Some p ->
if p = stolen then fail "the inherited path was bound: %S" p;
if not (await (fun () -> Sys.file_exists p)) then
fail "nothing was bound at the announced %S" p
else ignore (send p aso);
let reaped =
await ~ms:5000 (fun () ->
match Unix.waitpid [ Unix.WNOHANG ] spid with
| 0, _ -> false
| _ -> true)
in
if not reaped then begin
(try Unix.kill spid Sys.sigkill with Unix.Unix_error _ -> ());
fail "the program with an inherited variable never finished"
end);
(match In_channel.with_open_bin stolen In_channel.input_all with
| "the session's" -> ()
| _ -> fail "the inherited FLAN_AGENT_SOCKET's file was replaced"
| exception Sys_error _ ->
fail "the inherited FLAN_AGENT_SOCKET's file was removed");
(* Nobody's pid. *)
inherited ~shape:"an owner that is not this process" ~owner:"1";
(* A merged build's owner is the program itself, so a process the program
starts has the owner as its parent; with no FLAN_DEV_PARENT naming it,
that is not the --two-process shape and the socket is not its. Here the
test binary stands in for the program. *)
inherited ~shape:"a child of a merged program"
~owner:(string_of_int (Unix.getpid ()));
List.iter (fun f -> try Sys.remove f with Sys_error _ -> ())
[ aexe; aso; aout; aerr; bad; stolen; serr; tmp "stolen.out" ];
[ aexe; aso; aout; aerr; bad ];
(* ── No (agent/start) at all ────────────────────────────────────── *)

View File

@ -6784,7 +6784,38 @@ let () =
(%s)" backend
(Option.value ~default:"" (Wire.string_field r "value"))
(said r)
end
end;
(* A function value an expression makes has its code in that
expression's module — a lambda's body, or the wrapper a named
function is handed out through — so that module stays mapped.
Later expressions are mapped between the store and the call,
where an unloaded one would have been. *)
let defd code =
let r =
request c
(Printf.sprintf
"(:op \"eval\" :code %s \
:file \"programs/dev-dyn-global.flan\")" (Wire.quote code))
in
if status r <> "ok" then fail "--%s: %s: %s" backend code (said r)
in
defd "(defonce kept (Option (Fn [i64] i64)))";
defd "(defn twice [x i64] i64 (* x 2))";
let call_kept want what =
for i = 1 to 3 do
ignore (ev (Printf.sprintf "(do (println \"pad %d\") %d)" i i))
done;
let r = ev "(match kept (Some f) (f 1) (None) -1)" in
if Wire.string_field r "value" <> Some want then
fail "--%s: %s kept by an unloaded expression answered %S \
(%s)" backend what
(Option.value ~default:"" (Wire.string_field r "value"))
(said r)
in
ignore (ev "(do (set kept (Some (fn [x] (+ x 7)))) 0)");
call_kept "8" "a lambda";
ignore (ev "(do (set kept (Some twice)) 0)");
call_kept "2" "a named function"
end;
ignore (request c "(:op \"close\")");
(try Unix.close c with Unix.Unix_error _ -> ());
@ -8958,6 +8989,23 @@ let () =
"(:op \"eval-expr\" :code %s :file \"programs/dev-own-break.flan\")"
(Wire.quote code))
in
(* An expression stopped in a break and resumed by a restart finishes
after the restart's reply, and here it finishes after the next
expression has been sent: its value must not answer for that one. *)
let r =
ev "(restart-case (do (error (Late {})) 0) \
(slow [] (do (usleep 500000) 5)))"
in
if status r <> "error" then
fail "resumed value: the first expression did not stop: %s" (said r)
else begin
let r = request c "(:op \"restart\" :name \"slow\")" in
if status r <> "ok" then fail "resumed value: restart: %s" (said r);
let r = ev "(do (usleep 300000) 23)" in
if Wire.string_field r "value" <> Some "23" then
fail "resumed value: the next expression answered %S (%s)"
(Option.value ~default:"" (Wire.string_field r "value")) (said r)
end;
let r = ev "(do (set go 1) 0)" in
if status r <> "ok" then fail "own break: setting go: %s" (said r)
else begin

View File

@ -214,8 +214,19 @@ typedef struct {
void *handle;
int stopped_only;
int32_t at_stop;
uint32_t call_id; /* its place among jobs with a call; 0 if none */
} job;
/* Which evaluated expression the result buffer holds. Every job with a call
* is numbered as the listener queues it, and a call that returns records its
* number, so a daemon waiting for its own expression's value is not answered
* by an earlier expression that a restart resumed and that published after
* the new one was sent. The highest wins: an expression run inside another's
* break returns first, and the outer one only resumes on a later request.
* The [calls] verb answers both counts. */
static _Atomic uint32_t calls_queued;
static _Atomic uint32_t calls_valued;
/* Said once, in one place, and shipped to the daemon over [refusals] rather
* than written down again at the other end. A refusal is a sentence naming
* what actually happened, and the thing that actually happened is not "the
@ -1335,6 +1346,8 @@ int32_t flan_agent_poll(void) {
if (sigsetjmp(escape, 1) == 0) {
eval_escape = &escape;
j.call();
if (j.call_id > atomic_load(&calls_valued))
atomic_store(&calls_valued, j.call_id);
} else {
if (flan_condition_stacks_restore) flan_condition_stacks_restore(mh, mr, md);
if (flan_dev_frames_restore) flan_dev_frames_restore(mf);
@ -1920,6 +1933,14 @@ static void handle_line(char *line, sink *o) {
* editor polls this without knowing the state already. */
/* After the number, whose code stopped: "eval" when the thread was inside
* an evaluated thunk, "program" when it was in the program's own code. */
if (strcmp(line, "calls") == 0) {
char hdr[48];
int k = snprintf(hdr, sizeof hdr, "%u %u\n",
(unsigned)atomic_load(&calls_queued),
(unsigned)atomic_load(&calls_valued));
if (k > 0) emit(o, hdr, (size_t)k);
return;
}
if (strcmp(line, "stop") == 0) {
snapshot *s = (atomic_load(&depth) > 0) ? snap_top() : NULL;
char hdr[32];
@ -2199,7 +2220,9 @@ static void handle_line(char *line, sink *o) {
* is the failure being fixed. */
if (!publish((job){ .install = f, .call = c,
.handle = transient == NULL ? NULL : h,
.stopped_only = stopped_only, .at_stop = at_stop }))
.stopped_only = stopped_only, .at_stop = at_stop,
.call_id = c == NULL ? 0
: atomic_fetch_add(&calls_queued, 1) + 1 }))
fprintf(stderr, "flan: reload queue full after it was checked\n");
return;
}
@ -2498,8 +2521,17 @@ static const char *daemon_socket(void) {
return NULL;
pid = strtol(own, &end, 10);
if (end == own || *end != '\0' || pid <= 0) return NULL;
if (pid != (long)getpid() && pid != (long)getppid()) return NULL;
return env;
if (pid == (long)getpid()) return env;
/* The parent only under --two-process, which is the one shape that sets
* FLAN_DEV_PARENT, and to the same pid. In a merged build the owner is the
* program itself, so a process it starts has the owner as its parent and
* must not take the socket. */
{
const char *par = getenv("FLAN_DEV_PARENT");
if (par != NULL && strcmp(par, own) == 0 && pid == (long)getppid())
return env;
}
return NULL;
}
/* [path] is a Flan string: ptr and len, not NUL-terminated.