|
6 | 6 | [clojisr.v1.impl.rserve.call :as call] |
7 | 7 | [clojisr.v1.impl.rserve.packages :as packages] |
8 | 8 | [clojisr.v1.impl.rserve.printing :as printing] |
9 | | - [clojure.core.async :as async] |
10 | 9 | [clojure.tools.logging.readable :as log] |
11 | 10 | [clojisr.v1.util :refer [exception-cause get-free-port]]) |
12 | 11 | (:import [org.rosuda.REngine.Rserve RConnection] |
13 | 12 | [java.io BufferedReader])) |
14 | 13 |
|
15 | | -(def defaults |
16 | | - (atom |
17 | | - {:host "localhost" |
18 | | - :spawn-rserve? true})) |
| 14 | +(def defaults (atom {:host "localhost" |
| 15 | + :spawn-rserve? true})) |
19 | 16 |
|
20 | 17 | (defn close! [{:keys [^RConnection r-connection rserve]}] |
21 | 18 | (when r-connection |
22 | 19 | (.close r-connection)) |
23 | 20 | (when rserve |
24 | | - (proc/close rserve)) |
25 | | - nil) |
| 21 | + (proc/close rserve))) |
26 | 22 |
|
27 | 23 | ;; Session is valid when there is connection, also when we have rserve, process should be active |
28 | 24 | ;; if something is not true, ensure cleaning the rest and close the session |
29 | 25 | (defn active?-or-close! [{:keys [^RConnection r-connection rserve] |
30 | | - :as sess}] |
| 26 | + :as session}] |
31 | 27 | (let [state (and r-connection |
32 | 28 | (.isConnected r-connection) |
33 | 29 | (if-not rserve |
34 | 30 | true |
35 | 31 | (proc/alive? rserve)))] |
36 | | - (or state (close! sess)))) |
| 32 | + (or state (close! session)))) |
37 | 33 |
|
38 | 34 | (defrecord RserveSession [id |
39 | 35 | session-args |
|
44 | 40 | (close! session)) |
45 | 41 | (closed? [session] |
46 | 42 | (not (active?-or-close! session))) |
47 | | - (id [_session] |
48 | | - id) |
49 | | - (session-args [_session] |
50 | | - session-args) |
51 | | - (desc [_session] |
52 | | - session-args) |
| 43 | + (id [_session] id) |
| 44 | + (session-args [_session] session-args) |
| 45 | + (desc [_session] session-args) |
53 | 46 | (eval-r->java [session code] |
54 | 47 | (log/debug [::eval-r->java {:code code |
55 | 48 | :session-args (:session-args session)}]) |
|
58 | 51 | ;; Unlike (.assign r-connection ...), the following approach |
59 | 52 | ;; allows for a varname like "abc$xyz". |
60 | 53 | (locking r-connection |
61 | | - (.eval |
62 | | - r-connection |
63 | | - (call/assignment varname java-obj) |
64 | | - nil |
65 | | - true))) |
| 54 | + (.eval r-connection (call/assignment varname java-obj) nil true))) |
66 | 55 | (print-to-string [session r-obj] |
67 | 56 | (printing/print-to-string session r-obj)) |
68 | 57 | (package-symbol->r-symbol-names [session package-symbol] |
69 | | - (packages/package-symbol->r-symbol-names |
70 | | - session package-symbol)) |
| 58 | + (packages/package-symbol->r-symbol-names session package-symbol)) |
71 | 59 |
|
72 | 60 | iprot/Engine |
73 | 61 | (->nil [_] (rexp/->rexp-nil)) |
|
82 | 70 | (->named-list [_ ks vs] (rexp/->rexp-named-list ks vs)) |
83 | 71 | (native? [_ x] (rexp/rexp? x))) |
84 | 72 |
|
85 | | -(defn rserve-print-loop [{:keys [rserve] |
86 | | - :as session}] |
87 | | - (log/info [::rserve-print-loop {:action :started |
88 | | - :session-args (:session-args session)}]) |
89 | | - (async/go-loop [] |
90 | | - (doseq [^BufferedReader reader |
91 | | - (-> rserve |
92 | | - ((juxt :out :err)))] |
93 | | - (loop [] |
94 | | - (when (.ready reader) |
95 | | - (let [line (.readLine reader)] |
96 | | - (when-not |
97 | | - (re-find |
98 | | - ;; Just avoidingg this confusing message. |
99 | | - #"(This session will block until Rserve is shut down)" line) |
100 | | - (println line))) |
101 | | - (recur)))) |
| 73 | +(defn print-loop-task |
| 74 | + [{:keys [rserve] :as session}] |
| 75 | + (fn [] (doseq [[k output-stream] [[:out *out*] [:err *err*]]] |
| 76 | + (let [^BufferedReader reader (-> rserve k)] |
| 77 | + (binding [*out* output-stream] |
| 78 | + (loop [] |
| 79 | + (when (.ready reader) |
| 80 | + (let [line (.readLine reader)] |
| 81 | + (when-not |
| 82 | + (re-find |
| 83 | + ;; Just avoidingg this confusing message. |
| 84 | + #"(This session will block until Rserve is shut down)" line) |
| 85 | + (println line))) |
| 86 | + (recur)))))) |
102 | 87 | (Thread/sleep 100) |
103 | 88 | (if (not (prot/closed? session)) |
104 | 89 | (recur) |
105 | 90 | (log/info [::rserve-print-loop {:action :stopped |
106 | 91 | :session-args (:session-args session)}])))) |
107 | 92 |
|
| 93 | +(defn rserve-print-loop [session] |
| 94 | + (log/info [::rserve-print-loop {:action :started |
| 95 | + :session-args (:session-args session)}]) |
| 96 | + (.start (Thread. (print-loop-task session)))) |
| 97 | + |
| 98 | +(defn make |
| 99 | + "Creates RServe session. |
108 | 100 |
|
109 | | -(defn make [id session-args] |
110 | | - (let [{:keys [host port spawn-rserve? init-r]} (merge @defaults |
111 | | - session-args) |
| 101 | + Process is spawned (optionally), then connection is established." |
| 102 | + [id session-args] |
| 103 | + (let [{:keys [host port spawn-rserve? init-r] :as args} (merge @defaults session-args) |
112 | 104 | port (or port (get-free-port)) |
113 | 105 | rserve (when spawn-rserve? |
114 | 106 | (proc/start-rserve port init-r))] |
|
119 | 111 | (if (or (zero? attempts) |
120 | 112 | (proc/alive? rserve)) |
121 | 113 | (when-not (proc/alive? rserve) |
122 | | - (throw (Exception. "Can't create RServe process."))) |
| 114 | + (throw (ex-info "Can't create RServe process." args))) |
123 | 115 | (do |
124 | 116 | (log/warn [::rserve-spawn {:message "Rserve is not alive yet, waiting 0.5s"}]) |
125 | 117 | (recur (dec attempts)))))) |
126 | 118 |
|
127 | 119 | (let [conn (loop [attempts (int 1)] ;; try 5 times to connect |
128 | 120 | (when (> attempts 5) ;; throw an Exception when can't connect |
129 | | - (throw (Exception. "Can't connect to RServe, please check host/port settings."))) |
| 121 | + (throw (ex-info "Can't connect to RServe, please check host/port settings." args))) |
130 | 122 | (Thread/sleep (* attempts 200)) |
131 | 123 | (let [conn (try |
132 | 124 | (RConnection. host port) |
|
142 | 134 | session (->RserveSession id |
143 | 135 | session-args |
144 | 136 | conn |
145 | | - rserve)] |
| 137 | + rserve)] |
146 | 138 | (when rserve (rserve-print-loop session)) |
147 | 139 | session))) |
0 commit comments