Author SHA1 Message Date
Roman Volosovskyi cbb6962bfa task preprocessing 2016-04-21 12:19:57 +03:00
Roman Volosovskyi edefca55b4 ranking/revision first attempt 2016-04-20 17:45:13 +03:00
Roman Volosovskyi ae67c2aa32 quests LI 2016-04-14 17:21:16 +03:00
Roman Volosovskyi 5a20fd4b51 questions LI 2016-04-14 17:02:19 +03:00
Roman Volosovskyi d947c82ad3 goals LI tasklinks, update-quests-budgets, update-belief-budget 2016-04-14 16:37:44 +03:00
Roman Volosovskyi 5bdd2c4362 beliefs LI update-questions-budgets 2016-04-14 14:38:03 +03:00
Roman Volosovskyi d6975b6486 beliefs LI update-goal-budget 2016-04-14 14:29:44 +03:00
Roman Volosovskyi b1a49b7edd test new tasklinks 2016-04-14 13:12:10 +03:00
Roman Volosovskyi 1cf7c7b4d3 beliefs LI check-questions 2016-04-14 12:45:45 +03:00
Roman Volosovskyi 5a6a1ead36 LI jobs refactoring 2016-04-14 11:47:54 +03:00
Roman Volosovskyi 99e32287d1 goals/beliefs LI refactoring 2016-04-14 09:54:18 +03:00
Roman Volosovskyi b87ea44138 quests LI proto 2016-04-13 12:03:48 +03:00
Roman Volosovskyi 6677f41a12 questions LI proto 2016-04-13 11:31:57 +03:00
Roman Volosovskyi 181e690a65 goals LI proto 2016-04-13 10:02:08 +03:00
Roman Volosovskyi b204657325 beliefs LI proto 2016-04-12 18:10:56 +03:00
Roman Volosovskyi 333ef9a145 tableless memory 2016-04-08 18:31:45 +03:00
Roman Volosovskyi d88c402696 local inference mock 2016-04-07 18:58:39 +03:00
Roman Volosovskyi bfbbb14084 travis-ci upgrade 2016-04-07 12:50:27 +03:00
Roman Volosovskyi a6a3777915 selection of data from redis 2016-04-07 11:42:43 +03:00
Roman Volosovskyi e0cef19049 flow conditions & injections 2016-04-07 09:31:40 +03:00
Roman Volosovskyi 5b221cdfb8 jdk8 2016-04-06 17:17:43 +03:00
Roman Volosovskyi aeb037b8ce at least kafka works 2016-04-06 17:07:06 +03:00
Roman Volosovskyi ee666d4012 api updates 2016-04-04 15:14:10 +03:00
32 changed files with 1794 additions and 535 deletions
+1
View File
@@ -15,3 +15,4 @@ pom.xml.asc
\#*\#
.\#*
/src/nal/experiments.clj
*.log
+6
View File
@@ -6,3 +6,9 @@ notifications:
secure: eI5hj2PABtAUkpeYwiAHcFMa5pNvHFuU87uqe4qBtv0upPtu22/dLpw8l6EhhcMZlW6C8D2irrj6W5C1lujCy7u3Vpu1FeBPZbqzapyU1Hpe8YgSRdCY1dbygawZp/Av/HWzpi0aTket83F5/0dvTZK0fQTBRmAfyx17DEJscYEBH61dKVqdwBy5OVYvhQ1+QtEnt1TRJcT0AA8von9lanzx5/mWdQ6O+3BrKURWG6vvQUYEMZH43NU6lNGNV+PBGV1PSqZqbzCh/6C6/d6HauQdewr4Oubez90OquyCDpCZDoVG8eCAsZtZt+tzfJ5gjtz8+HA6GQeQB6QMLAp89m1bPwSZCqdnmQp9S8TNgm8Io3jREBh+5JIdpDXwukKT1kMPFrdDiPTEhHNJojyK3/BERGyrhf8azow+brPq0EIM9Fi/SkGa0gb9mUXY/BZ1MF7ulxxbzpLxHJYlT9QVlV9Q09/uhH4tIh0kPQhe+ntkMT1uLTfDur7CZakB5270ibreRSeK4RKdpy5SjaIN70hSwyrWoZHiED6aKEJJSO1t9Ve1jTDyWbcUZYewtBi+APcqAraja9NIlFApjBRO1jUneS4/BzNqxWZAFS7MOwMrk1xFft47hMIHoNYd7IgyB7TWSxhloNTwMyOWkdOhnPIqqCvds6M0yQCDKNBiVeM=
services:
- redis-server
jdk:
- oraclejdk8
sudo: required
dist: precise
group: edge
+9 -3
View File
@@ -10,9 +10,13 @@
[org.clojure/data.priority-map "0.0.7"]
[org.clojure/core.match "0.3.0-alpha4"]
[org.clojure/core.unify "0.5.5"]
[org.clojure/core.async "0.2.374"]
[org.onyxplatform/onyx "0.9.0"]
[org.onyxplatform/onyx-kafka "0.9.0.1"]
[com.taoensso/carmine "2.12.2"]
[mount "0.1.10"]]
[clj-kafka "0.3.4" :exclusions [org.apache.zookeeper/zookeeper zookeeper-clj]]
[mount "0.1.10"]
[aero "0.2.1"]
[com.taoensso/timbre "4.1.4"]]
:main ^:skip-aot narjure.core
:plugins [[lein-cloverage "1.0.6"]
[jonase/eastwood "0.2.3"]
@@ -22,4 +26,6 @@
:target-path "target/%s"
:repl-options {:init-ns narjure.repl
:nrepl-middleware [narjure.repl/narsese-handler]}
:profiles {:uberjar {:aot :all}})
:profiles {:uberjar {:aot :all}
:test {:dependencies [[com.stuartsierra/component "0.3.1"]]}
:dev {:dependencies [[com.stuartsierra/component "0.3.1"]]}})
+29
View File
@@ -0,0 +1,29 @@
{:env-config
{:onyx/tenancy-id #env [ONYX_ID "testcluster"]
:onyx.bookkeeper/server? true
:onyx.bookkeeper/local-quorum? #cond {:default false
:test true}
:onyx.bookkeeper/delete-server-data? true
:onyx.bookkeeper/local-quorum-ports [3196 3197 3198]
:onyx.bookkeeper/port 3196
:zookeeper/address #cond {:default #env [ZOOKEEPER "zookeeper"]
:test "127.0.0.1:2188"}
:zookeeper/server? #cond {:default false
:test true}
:zookeeper.server/port 2188}
:peer-config
{:onyx/tenancy-id #env [ONYX_ID "testcluster"]
:zookeeper/address #cond {:default #env [ZOOKEEPER "zookeeper"]
:test "127.0.0.1:2188"}
:onyx.peer/job-scheduler :onyx.job-scheduler/greedy
:onyx.peer/zookeeper-timeout 60000
:onyx.messaging/allow-short-circuit? #cond {:default true
:test false}
:onyx.messaging/impl :aeron
:onyx.messaging/bind-addr #env [BIND_ADDR "localhost"]
:onyx.messaging/peer-port 40200
:onyx.messaging.aeron/embedded-driver? #cond {:default false
:test true}
:onyx.messaging/backpressure-strategy :high-restart-latency}
:redis-config
{:uri "redis://127.0.0.1:6379"}}
+1 -7
View File
@@ -1,7 +1,6 @@
(ns nal.core
(:require [nal.deriver.truth :as t]
[nal.deriver :refer [generate-conclusions]]
[nal.rules :as r]))
[nal.deriver :refer [generate-conclusions]]))
(defn choice [[f1 c1] [f2 c2]]
(if (>= c1 c2) [f1 c1] [f2 c2]))
@@ -11,8 +10,3 @@
(generate-conclusions (rules task-type) task belief))
(def revision t/revision)
(comment
:shift-occurrence-forward ;pre
:shift-occurrence-backward ;pre
:linkage-temporal)
+1 -1
View File
@@ -23,7 +23,7 @@
`n/reduce-int-dif `n/reduce-and `n/reduce-ext-dif `n/reduce-image
`n/reduce-int-inter `n/reduce-neg `n/reduce-or `nil? `not `or `abs
`implications-and-equivalences `get-terms `empty? `intersection
`n/reduce-seq-conj `clojure.core.match/match})
`n/reduce-seq-conj})
(defn operators->placeholders
[statement]
+4
View File
@@ -3,6 +3,7 @@
;https://github.com/opennars/opennars/blob/6611ee7f0b1428676b01ae4a382241a77ae3a346/nars_logic/src/main/java/nars/nal/meta/BeliefFunction.java
;https://github.com/opennars/opennars/blob/7b27dacec4cdbe77ca03d89296323d49875ac213/nars_logic/src/main/java/nars/truth/TruthFunctions.java
;https://github.com/opennars/opennars/blob/b3d97d1543b70dac604aff75eca7fe3d38f1e823/nars_logic/src/main/java/nars/truth/Truth.java#L87
(defn t-and
(^double [^double a ^double b] (* a b))
@@ -131,6 +132,9 @@
[t _]
(analogy t [1.0 d/judgement-confidence]))
(defn expectation [[f c]]
(+ (* c (- f 0.5)) 0.5))
(def tvtypes
{:t/structural-deduction structual-abduction
:t/struct-int structual-intersection
-31
View File
@@ -1,31 +0,0 @@
(ns narjure.control.buffer
(:require [clojure.core.async.impl.protocols :as impl])
(:import [java.util LinkedList]
[clojure.lang Fn Counted]))
(deftype PanickingSlidingBuffer
[^LinkedList buf ^long n ^Fn warning-callback ^long warning-n]
impl/UnblockingBuffer
impl/Buffer
(full? [this]
false)
(remove! [this]
(.removeLast buf))
(add!* [this itm]
(let [size (.size buf)]
(when (>= size warning-n)
(warning-callback size)
(when (= size n)
(impl/remove! this))))
(.addFirst buf itm)
this)
(close-buf! [this])
Counted
(count [this]
(.size buf)))
(defn panicking-sliding-buffer
([n callback]
(panicking-sliding-buffer n callback (Math/round (* 0.9 n))))
([n callback warning-n]
(PanickingSlidingBuffer. (LinkedList.) n callback warning-n)))
-110
View File
@@ -1,110 +0,0 @@
(ns narjure.control.flow
(:require [clojure.core.async :refer [go-loop <! >! chan]]
[clojure.set :as set]))
(defn check-element-in-map
"Checks if elements exist in map, if not
assocs elements to map with default value."
([s m] (check-element-in-map s 0 m))
([s default m]
(reduce (fn [ac k]
(if (ac k)
ac
(assoc ac k default)))
m s)))
(defn all
"Returns set aff all functions from workflow."
[wf]
(set (flatten wf)))
(defn kw->fn
"Transform function's keyword to var."
[kw]
(->> (str kw)
(drop 1)
(apply str)
symbol
resolve))
(defn vertex
"Creates vertex of flow graph. Arguments:
- functions: collection of collections, where first element is function
and second (optional) is output port
- inputs: ports (edges) that should be listened by vertex
- p: number of parallelism"
[functions inputs p]
(doseq [in inputs
_ (range (* p (count functions)))]
(go-loop []
(when-let [val (<! in)]
(doseq [[f out] functions]
(let [results (f val)]
(when out
(doseq [res (if (map? results)
[results]
results)]
(>! out res)))))
(recur)))))
(defn check-output
"Creates output port for function if it is necessary."
[buffer [function output-cnt]]
[function (when (pos? output-cnt) (chan buffer))])
(defn fn-outputs
"Generates map where keys are functions and values are ports which
will be used to send result of execution of functions."
[workflow buffer]
(->> (group-by first workflow)
(map (fn [[n t]] [n (count t)]))
(into {})
(check-element-in-map (all workflow))
(map #(check-output buffer %))
(into {})))
(defn fn-inputs [workflow]
"Groups functions to identify vertexes and edges that they should listen.
Returns map {vertexes edges ...}
in: [[:a :b]
[:a :c]
[:c :d]
[:b :d]]
out: {[:c :b] [:a], [:d] [:b :c]}"
(->> workflow
(reduce (fn [ac [k v]]
(update ac k conj v))
{})
(reduce (fn [ac [k v]]
(update ac v conj k))
{})))
(defn generate-flow
"Generates flow of functions which is discribed by pairs of functions,
where result of fisrt function will be sent to input of the second function.
Optionally map of configuration params can be passed.
Possible configs:
- parallelism: map where keys are functions and values
are parallelization numbers
- default-p: defaulp parallelization number
- buffer: capacity of fixed buffer for all channels"
;TODO configuration for custom buffers
([workflow] (generate-flow workflow {}))
([workflow {:keys [parallelism default-p buffer]
:or {parallelism {}
default-p 1
buffer 100}}]
(let [in (chan buffer)
outputs (assoc (fn-outputs workflow buffer) :in in)
inputs (fn-inputs workflow)
all-inputs (set (mapcat key inputs))
input-tasks (set/difference (all workflow) all-inputs)
it2 (assoc inputs input-tasks [:in])]
(doseq [[tasks from] it2]
(vertex
(map (fn [f] [(kw->fn f) (outputs f)]) tasks)
(map outputs from)
(apply max (map #(parallelism % default-p) tasks))))
in)))
+108 -49
View File
@@ -1,53 +1,115 @@
(ns narjure.control.general-inference
(:require [narjure.system :refer [memory inference]]
[narjure.control.flow :as f]
[narjure.memory.api :as m]))
(:require
[onyx.job :refer [add-task]]
[onyx.plugin.kafka]
[onyx.tasks.kafka :refer [consumer producer]]
[narjure.system :refer [inference]]
[taoensso.timbre :refer [info]]
[narjure.memory.api :as m]
[narjure.control.utils :refer [leaf-function intermediate-function
inject-memory kafka-consumer
kafka-producer]]))
;; select-concepts
;; | | | |
;; v v v ... v
;; select-task-link --> update-tasklink-budget
;; |
;; v
;; select-term-link --> update-termlink-budget
;; |
;; v
;; do-inference
(def workflow
[[::select-concepts ::select-task-link]
[[:select-active-concepts :select-task-link]
[:select-task-link :update-tasklink-budget]
[:select-task-link :select-term-link]
[::select-task-link ::update-tasklink-budget]
[::select-task-link ::select-term-link]
[:select-term-link :update-termlink-budget]
[:select-term-link :do-general-inference]
[:do-general-inference :write-tasks]
[::select-term-link ::update-termlink-budget]
[::select-term-link ::do-inference]])
[:select-term-link :choose-answer]
[:choose-answer :write-answer]])
(def flow-conditions
[{:flow/from :select-term-link
:flow/to [:do-general-inference]
:flow/predicate [:not ::question-with-query-var?]}
{:flow/from :select-term-link
:flow/to [:choose-answer]
:flow/predicate ::question-with-query-var?}])
(defn select-concept [_]
(m/pull-activated-concepts memory))
(defn question-with-query-var? [_ _ segment _]
(if (:question-with-query-var segment) true false))
(defn select-task-link [concept]
(let [{task-id :task} (m/select-tasklink memory concept)
task (m/task memory task-id)]
{:task task
:concept concept}))
(defn catalog [batch-settings]
(mapv #(merge % batch-settings)
[(intermediate-function :select-task-link ::select-task-link)
(leaf-function :update-tasklink-budget ::update-tasklink-budget batch-settings)
(intermediate-function :select-term-link ::select-term-link)
(leaf-function :update-termlink-budget ::update-termlink-budget batch-settings)
(intermediate-function :do-general-inference ::do-general-inference)
(intermediate-function :choose-answer ::choose-answer)]))
;(defn update-concept-budget [data] data)
(def tasks-with-redis-access
(disj (set (flatten workflow))
:select-active-concepts
:write-tasks
:write-answer
:do-general-inference))
(defn select-term-link [{:keys [concept task] :as data}]
(let [{linked-concept :concept} (m/select-termlink memory concept)
(def lifecycles
(concat
(map #(hash-map :lifecycle/task %
:lifecycle/calls :narjure.control.utils/memory-call)
tasks-with-redis-access)
[{:lifecycle/task :do-general-inference
:lifecycle/calls ::inference-call}]))
(defn inject-inference
[{params :onyx.core/params} _]
{:onyx.core/params (conj params inference)})
(def inference-call
{:lifecycle/before-task-start inject-inference})
(defn build-job
[zk-address concepts-topic tasks-topic ansers-topic batch-size batch-timeout]
(let [batch-settings {:onyx/batch-size batch-size
:onyx/batch-timeout batch-timeout}
base-job (merge {:workflow workflow
:catalog (catalog batch-settings)
:lifecycles lifecycles
:windows []
:triggers []
:flow-conditions flow-conditions
:task-scheduler :onyx.task-scheduler/balanced})]
(-> base-job
(add-task (kafka-consumer :select-active-concepts concepts-topic
zk-address batch-settings 1))
(add-task (kafka-producer :write-tasks tasks-topic zk-address
batch-settings))
(add-task (kafka-producer :write-answer ansers-topic zk-address
batch-settings)))))
;-------------------------------------------------------------------------------
;functions
(defn select-task-link [mem {:keys [concept] :as segment}]
(let [task (m/select-task mem concept)]
(merge segment {:task task
:concept concept})))
(defn update-tasklink-budget [mem segment]
(info :update-tasklink-budget))
(defn select-term-link
[mem {:keys [concept task] :as segment}]
(let [{linked-concept :concept} (m/select-termlink mem concept)
occurrence (:occurrence task)
truth (m/select-truth memory linked-concept occurrence)]
(assoc data :truth truth
:term (m/term memory linked-concept))))
;(defn update-tasklink-budget [data] data)
;(defn update-termlink-budget [data] data)
belief (m/select-belief mem linked-concept occurrence)]
(info :belief belief)
(assoc segment :belief belief)))
(defn do-inference [{:keys [task truth term]}]
(let [{statement :term
:keys [frequency confidence plausibility
desirability task-type occurrence]}
task
(defn update-termlink-budget [mem segment]
(info :update-termlink-budget))
(defn do-general-inference
[inference {:keys [task belief]}]
(let [belief (:task belief)
{:keys [frequency confidence plausibility statement
desirability task-type occurrence]} (:task task)
task {:statement statement
:desire [plausibility desirability]
@@ -55,15 +117,12 @@
:task-type task-type
:occurrence occurrence}
truth (assoc truth :statement term)
results (inference task truth)]
(doseq [res results]
(m/push-task memory res))))
{:keys [confidence frequency]} belief
belief (assoc belief :truth [frequency confidence])]
(info :inference task belief)
(mapv (fn [conclusion] {:message conclusion})
(inference task belief))))
(def default-parallelism {::do-inference 4})
(defn general-inference-flow
[{:keys [parallelism]
:or {parallelism default-parallelism}
:as options}]
(f/generate-flow workflow options))
(defn choose-answer [mem segment]
(info :choose-answer)
{:message {:some-answer true}})
-37
View File
@@ -1,37 +0,0 @@
(ns narjure.control.local-inference)
(def workflow
[[:read-task :answer-yn-question]
[:read-task :answer-general-question]
[:answer-yn-question :out-answers]
[:answer-general-question :out-answers]
[:read-task :find-related-concepts]
[:find-related-concepts :out-update-tasklinks]
[:find-related-concepts :out-check-tasklinks-capacity]
[:read-task :belief-revision]
[:belief-revision :out-update-beliefs]
[:read-task :goal-revision]
[:goal-revision :out-update-goals]])
(defn answer-yn-question [{:keys [task] :as segment}]
(println :yn)
{})
(defn answer-general-question [{:keys [task] :as segment}]
(println :general)
{})
(defn find-related-concepts [{:keys [task] :as segment}]
(println :rel)
{})
(defn belief-revision [{:keys [task] :as segment}]
(println :bel)
{})
(defn goal-revision [{:keys [task] :as segment}]
(println :goal)
{})
@@ -0,0 +1,129 @@
(ns narjure.control.local-inference.beliefs
(:require
[onyx.plugin.kafka]
[narjure.memory.api :as m]
[taoensso.timbre :refer [info]]
[narjure.control.local-inference.utils :as u]))
(def workflow
[[:read-beliefs :check-satisfaction]
[:check-satisfaction :revise-anticipations]
[:revise-anticipations :revise-beliefs]
[:revise-beliefs :check-questions]
[:check-questions :add-tasklink]
[:check-questions :add-revised-tasklink]
[:check-questions :update-questions-budgets]
[:check-questions :update-goal-budget]
[:check-questions :prepare-answers]
[:prepare-answers :write-answers]
[:check-questions :add-anticipation-tasklink]])
(def flow-conditions
[{:flow/from :check-questions
:flow/to [:add-revised-tasklink]
:flow/predicate :narjure.control.local-inference.utils/has-revised-task?}
{:flow/from :check-questions
:flow/to [:add-anticipation-tasklink]
:flow/predicate ::observable?}
{:flow/from :check-questions
:flow/to [:add-tasklink]
:flow/predicate :narjure.control.local-inference.utils/constantly-true}
{:flow/from :check-questions
:flow/to [:update-questions-budgets :prepare-answers]
:flow/predicate :narjure.control.local-inference.utils/has-answers?}
{:flow/from :check-questions
:flow/to [:update-goal-budget]
:flow/predicate ::satisfied-goal?}])
(def kafka-io #{:read-beliefs :write-answers})
(def leaves
#{:add-tasklink :add-revised-tasklink :add-anticipation-tasklink
:update-questions-budgets :update-goal-budget})
(def build-job
(u/get-job-builder (namespace ::k) workflow kafka-io leaves flow-conditions))
;-------------------------------------------------------------------------------
;predicates
(defn observable? [_ _ {:keys [revised-task]} _]
;TODO implement
false)
(defn satisfied-goal? [_ _ {:keys [satisfied-goal]} _]
satisfied-goal)
;-------------------------------------------------------------------------------
;functions
(defn check-satisfaction
[mem {{:keys [statement] :as task} :task :as segment}]
(if-let [goals (m/goals mem statement)]
(let [goal-tasklink (u/rank goals task)
satisfaction (u/satisfaction (:task goal-tasklink) task)]
(-> segment
(update-in [:task :durability] u/increase-durability satisfaction)
(assoc :satisfied-goal
{:goal goal-tasklink
:priority-diff (u/priority-diff goal-tasklink satisfaction)})))
segment))
;TODO implement
(defn revise-anticipations [_ segment]
segment)
(defn revise-beliefs
[mem {{:keys [statement] :as task} :task :as segment}]
(let [beliefs (m/beliefs mem statement)]
(if-let [revised-task (u/revise beliefs task)]
(assoc segment :revised-task (dissoc revised-task :id))
segment)))
(defn answer-question
[question-tasklink answer]
(let [satisfaction (u/q-satisfaction answer)]
{:question question-tasklink
:satisfaction satisfaction
:priority-diff (u/priority-diff question-tasklink satisfaction)}))
(defn check-questions
[mem {{:keys [statement] :as task} :task
revised-task :revised-task
:as segment}]
(if-let [questions (m/questions mem statement)]
(let [answer (or revised-task task)
results (mapv #(answer-question % answer) questions)
answer (reduce
(fn [answer {:keys [satisfaction]}]
(update answer :durability u/increase-durability satisfaction))
answer
results)]
(assoc segment
:answers results
(if revised-task :revised-task :task) answer))
segment))
(defn prepare-answers [_ {:keys [answers]}]
(mapv (fn [answer] {:message answer}) answers))
(defn add-tasklink
[mem {:keys [task]}]
(u/add-tasklink* mem (:statement task) task "belief"))
(defn add-revised-tasklink
[mem {:keys [revised-task]}]
(let [id (m/add-task mem revised-task)
revised-task' (assoc revised-task :id id)]
(u/add-tasklink* mem (:statement revised-task) revised-task' "belief")))
(defn update-questions-budgets
[mem {:keys [answers]}]
(mapv (fn [{:keys [question priority-diff]}]
(m/increment-value mem (question :id) :priority priority-diff))
answers))
(defn update-goal-budget
[mem {{:keys [goal priority-diff]} :satisfied-goal}]
(m/increment-value mem (goal :id) :priority priority-diff))
;TODO implement
(defn add-anticipation-tasklink [mem segment])
@@ -0,0 +1,135 @@
(ns narjure.control.local-inference.goals
(:require
[onyx.plugin.kafka]
[narjure.memory.api :as m]
[taoensso.timbre :refer [info]]
[narjure.control.local-inference.utils :as u]))
(def workflow
[[:read-goals :check-solution]
[:check-solution :revise-goals]
[:revise-goals :check-quests]
[:check-quests :check-operator]
[:check-operator :add-tasklink]
[:check-operator :add-revised-tasklink]
[:check-operator :update-quests-budgets]
[:check-operator :update-belief-budget]
[:check-operator :prepare-answers]
[:prepare-answers :write-answers]
[:check-operator :send-operator]])
(def flow-conditions
[{:flow/from :check-operator
:flow/to [:add-revised-tasklink]
:flow/predicate ::narjure.control.local-inference.utils/has-revised-task?}
{:flow/from :check-operator
:flow/to [:send-operator]
:flow/predicate ::has-operator-to-execute?}
{:flow/from :check-operator
:flow/to [:add-tasklink]
:flow/predicate :narjure.control.local-inference.utils/constantly-true}
{:flow/from :check-operator
:flow/to [:update-quests-budgets :prepare-answers]
:flow/predicate :narjure.control.local-inference.utils/has-answers?}
{:flow/from :check-operator
:flow/to [:update-belief-budget]
:flow/predicate ::has-solution?}])
(def io #{:read-goals :write-answers})
(def leaves
#{:add-tasklink :add-revised-tasklink :send-operator
:update-quests-budgets :update-belief-budget})
(def build-job
(u/get-job-builder (namespace ::k) workflow io leaves flow-conditions))
;-------------------------------------------------------------------------------
;predicates
(defn has-operator-to-execute? [_ _ _ _]
;TODO implement
false)
(defn has-solution? [_ _ {:keys [belief-solution]} _]
belief-solution)
;-------------------------------------------------------------------------------
;functions
(defn durability-diff
[{:keys [durability]} satisfaction]
(-> durability
(u/increase-durability satisfaction)
(- durability)))
(defn check-solution
[mem {{:keys [statement] :as task} :task :as segment}]
(if-let [beliefs (m/beliefs mem statement)]
(let [belief-tasklink (u/rank beliefs task)
satisfaction (u/satisfaction task (:task belief-tasklink))]
(-> segment
(update-in [:task :priority] u/reduce-priority satisfaction)
(assoc :belief-solution
{:belief belief-tasklink
:durability-diff (durability-diff belief-tasklink
satisfaction)})))
segment))
(defn revise-goals
[mem {{:keys [statement] :as task} :task :as segment}]
(let [goals (m/goals mem statement)]
(if-let [revised-task (u/revise goals task)]
(assoc segment :revised-task (dissoc revised-task :id))
segment)))
(defn answer-question
[quest-tasklink answer]
(let [satisfaction (u/q-satisfaction answer)]
{:quest quest-tasklink
:satisfaction satisfaction
:priority-diff (u/priority-diff quest-tasklink satisfaction)}))
(defn check-quests
[mem {{:keys [statement] :as task} :task
revised-task :revised-task
:as segment}]
(if-let [quests (m/quests mem statement)]
(let [answer (or revised-task task)
results (mapv #(answer-question % answer) quests)
answer (reduce
(fn [answer {:keys [satisfaction]}]
(update answer :durability u/increase-durability satisfaction))
answer
results)]
(assoc segment
:answers results
(if revised-task :revised-task :task) answer))
segment))
;TODO implement
(defn check-operator [_ segment]
segment)
(defn prepare-answers [_ {:keys [answers]}]
(mapv (fn [answer] {:message answer}) answers))
(defn add-tasklink
[mem {:keys [task]}]
(u/add-tasklink* mem (:statement task) task "goal"))
(defn add-revised-tasklink
[mem {:keys [revised-task]}]
(let [id (m/add-task mem revised-task)
revised-task' (assoc revised-task :id id)]
(u/add-tasklink* mem (:statement revised-task) revised-task' "goal")))
(defn update-quests-budgets
[mem {:keys [answers]}]
(mapv (fn [{:keys [quest priority-diff]}]
(m/increment-value mem (quest :id) :priority priority-diff))
answers))
(defn update-belief-budget
[mem {{:keys [belief durability-diff]} :belief-solution}]
(m/increment-value mem (belief :id) :durability durability-diff))
(defn send-operator [mem segment])
@@ -0,0 +1,57 @@
(ns narjure.control.local-inference.questions
(:require
[narjure.memory.api :as m]
[onyx.job :refer [add-task]]
[taoensso.timbre :refer [info]]
[narjure.control.local-inference.utils :as u]))
(def workflow
[[:read-questions :check-answer]
[:check-answer :add-tasklink]
[:check-answer :update-belief-budget]
[:check-answer :prepeare-answer]
[:prepeare-answer :write-answer]])
(def flow-conditions
[{:flow/from :check-answer
:flow/to [:update-belief-budget :prepeare-answer]
:flow/predicate :narjure.control.local-inference.utils/has-answer?}
{:flow/from :check-answer
:flow/to [:add-tasklink]
:flow/predicate :narjure.control.local-inference.utils/constantly-true}])
(def kafka-io #{:read-questions :write-answer})
(def leaves #{:add-tasklink :update-belief-budget})
(def build-job
(u/get-job-builder (namespace ::k) workflow kafka-io leaves flow-conditions))
;-------------------------------------------------------------------------------
;functions
(defn check-answer
[mem {{:keys [statement] :as task} :task
:as segment}]
(if-let [beliefs (m/beliefs mem statement)]
(let [belief-tasklink (u/rank beliefs task)
satisfaction (u/q-satisfaction (:task belief-tasklink))
task' (update task :priority u/reduce-priority satisfaction)]
(assoc segment
:task task'
:answer {:belief belief-tasklink
:question task'
:durability-diff (u/durability-diff belief-tasklink
satisfaction)}))
segment))
(defn add-tasklink
[mem {:keys [task]}]
(u/add-tasklink* mem (:statement task) task "question"))
(defn prepeare-answer [_ {:keys [answer]}]
;TODO check if input question
{:message answer})
(defn update-belief-budget
[mem {{:keys [belief durability-diff]} :answer}]
(m/increment-value mem (belief :id) :durability durability-diff))
@@ -0,0 +1,57 @@
(ns narjure.control.local-inference.quests
(:require
[narjure.memory.api :as m]
[onyx.job :refer [add-task]]
[taoensso.timbre :refer [info]]
[narjure.control.local-inference.utils :as u]))
(def workflow
[[:read-quests :check-answer]
[:check-answer :add-tasklink]
[:check-answer :update-goal-budget]
[:check-answer :prepeare-answer]
[:prepeare-answer :write-answer]])
(def flow-conditions
[{:flow/from :check-answer
:flow/to [:update-goal-budget :prepeare-answer]
:flow/predicate :narjure.control.local-inference.utils/has-answer?}
{:flow/from :check-answer
:flow/to [:add-tasklink]
:flow/predicate :narjure.control.local-inference.utils/constantly-true}])
(def kafka-io #{:read-quests :write-answer})
(def leaves #{:add-tasklink :update-goal-budget})
(def build-job
(u/get-job-builder (namespace ::k) workflow kafka-io leaves flow-conditions))
;-------------------------------------------------------------------------------
;functions
(defn check-answer
[mem {{:keys [statement] :as task} :task
:as segment}]
(if-let [goals (m/goals mem statement)]
(let [goal-tasklink (u/rank goals task)
satisfaction (u/q-satisfaction (:task goal-tasklink))
task' (update task :priority u/reduce-priority satisfaction)]
(assoc segment
:task task'
:answer {:goal goal-tasklink
:quest task'
:durability-diff (u/durability-diff goal-tasklink
satisfaction)}))
segment))
(defn add-tasklink
[mem {:keys [task]}]
(u/add-tasklink* mem (:statement task) task "quest"))
(defn prepeare-answer [_ {:keys [answer]}]
;TODO check if input question
{:message answer})
(defn update-goal-budget
[mem {{:keys [goal durability-diff]} :answer}]
(m/increment-value mem (goal :id) :durability durability-diff))
@@ -0,0 +1,197 @@
(ns narjure.control.local-inference.utils
(:require
[onyx.plugin.kafka]
[narjure.memory.api :as m]
[nal.deriver.truth :as t]
[onyx.job :refer [add-task]]
[narjure.control.utils :refer [leaf-function intermediate-function
kafka-consumer kafka-producer]]
[clojure.set :as s]))
(defn has-answers? [_ _ {:keys [answers]} _] (seq answers))
(defn has-answer? [_ _ {:keys [answer]} _] answer)
(defn has-revised-task? [_ _ {:keys [revised-task]} _] revised-task)
(def constantly-true (constantly true))
(defn add-tasklink* [mem concept task type]
(let [link (-> (select-keys task [:priority :durability])
(assoc :task (:id task)))
id (m/add-tasklink mem concept link (or type "task"))]
id))
(defn get-frequency [{:keys [frequency plausibility]}]
[(if frequency :frequency :plausibility)
(or frequency plausibility)])
(defn get-confidence [{:keys [confidence desirability]}]
[(if confidence :confidence :desirability)
(or confidence desirability)])
(defn revision
[revised task2]
(let [[f1k f1] (get-frequency revised)
[_ f2] (get-frequency task2)
[c1k c1] (get-confidence revised)
[_ c2] (get-confidence task2)
c2' (- 1 c2)
c1' (- 1 c1)
c1-c2' (* c1 c2')
c2-c1' (* c2 c1')
x (+ c1-c2' c2-c1')
f (/ (+ (* f1 c1-c2') (* f2 c2-c1'))
x)
c (/ x (+ x (* c1' c2')))]
{f1k f
c1k c}))
(defn overlaps? [evidences1 evidences2]
(seq (s/intersection (set evidences1)
(set evidences2))))
(defn merge-evidences
[ev1 ev2]
(vec (s/union (set ev1) (set ev2))))
(defn eternalize
[{:keys [confidence desirability] :as task}]
(let [k (if confidence :confidence :desirability)
val (or confidence desirability)]
(assoc task k (t/w2c val))))
(defn project-confidence
[source-time target-time current-time confidence]
(let [kc (/ (Math/abs (- source-time target-time))
(+ (Math/abs (- source-time current-time))
(Math/abs (- target-time current-time))))]
(* (- 1 kc) confidence)))
(defn project
[{source :occurrence
eternal-source :eternal
:as projected}
{target :occurrence
eternal-target :eternal
creation :creation-time}]
(case [eternal-target eternal-source]
[true false] (eternalize projected)
[false false]
(if (= source target)
projected
(let [[k value] (get-confidence projected)
new-value (project-confidence source target creation value)]
(assoc projected k new-value)))
:else projected))
(defn revise [tasklinks task]
(let [revised (reduce (fn [{:keys [evidences] :as input}
{:keys [task]}]
(if (overlaps? evidences (:evidences task))
input
(let [
evidences (merge-evidences evidences
(:evidences task))
projected (project task input)
revised-tv (revision input projected)]
(-> input
(assoc :evidences evidences)
(merge revised-tv)))))
task tasklinks)]
(when-not (= revised task) revised)))
(defn rank [tasklinks input]
(reduce (fn [{ac-task :task :as ac}
{:keys [task] :as tasklink}]
(let [projected (project task input)
[_ p-confidence] (get-confidence projected)
[_ ac-confidence] (get-confidence ac-task)]
(if (and (not (nil? ac)) (> ac-confidence p-confidence))
ac
(assoc tasklink :task projected))))
nil
tasklinks))
(defn satisfaction* [d t]
(- 1 (Math/abs (- (t/expectation d) (t/expectation t)))))
(defn satisfaction
[{:keys [plausibility desirability]}
{:keys [frequency confidence]}]
(satisfaction* [plausibility desirability]
[frequency confidence]))
(defn q-satisfaction [{:keys [confidence desirability]}]
(- 1 (or confidence desirability)))
(defn reduce-priority [priority satisfaction]
(t/t-and priority (- 1 satisfaction)))
(defn priority-diff
[{:keys [priority]} satisfaction]
(-> priority
(reduce-priority satisfaction)
(- priority)))
(defn increase-durability [durability satisfaction]
;TODO define k somewhere
(let [k 1]
(t/t-or durability (* k satisfaction))))
(defn durability-diff
[{:keys [durability]} satisfaction]
(-> durability
(increase-durability satisfaction)
(- durability)))
(defn save-task
[mem {:keys [task] :as segment}]
(assoc-in segment [:task :id] (m/add-task mem task)))
(defn get-catalog-builder
[functions-ns all leaves io]
(fn [batch-settings]
(let [add-ns #(keyword (str functions-ns "/" (name %)))
intermediates (s/difference all leaves io)
leaf-tasks (map #(leaf-function % (add-ns %) batch-settings) leaves)
inter-tasks (map #(intermediate-function % (add-ns %)) intermediates)]
(mapv #(merge % batch-settings) (concat leaf-tasks inter-tasks)))))
(defn job-builder
[workflow catalog-f lifecycles flow-conditions io]
(fn [zk-address {:keys [input output]} batch-size batch-timeout]
(let [batch-settings {:onyx/batch-size batch-size
:onyx/batch-timeout batch-timeout}
base-job (merge {:workflow workflow
:catalog (catalog-f batch-settings)
:lifecycles lifecycles
:windows []
:triggers []
:flow-conditions flow-conditions
:task-scheduler :onyx.task-scheduler/balanced})]
(let [parameters (s/union (set (keys input)) (set (keys output)))]
(assert
(= parameters io)
(format "All consumers/producers must be specified, %s are missed."
(s/difference io parameters))))
(reduce
add-task
base-job
(concat
(map (fn [[task topic]]
(kafka-consumer task topic zk-address batch-settings 1)) input)
(map (fn [[task topic]]
(kafka-producer task topic zk-address batch-settings)) output))))))
(defn get-job-builder
[functions-ns workflow io leaves flow-conditions]
(let [all (set (flatten workflow))
catalog (get-catalog-builder functions-ns all leaves io)
tasks-with-redis-access (s/difference all io)
lifecycles
(mapv #(hash-map :lifecycle/task %
:lifecycle/calls :narjure.control.utils/memory-call)
tasks-with-redis-access)]
(job-builder workflow catalog lifecycles flow-conditions io)))
@@ -0,0 +1,79 @@
(ns narjure.control.task-preprocessing
(:require [narjure.control.local-inference.utils :as u]
[narjure.memory.api :as m]
[taoensso.timbre :refer [info]]
[clojure.set :as s]))
(def workflow
[[:read-task :save-task]
[:save-task :prepare-task]
[:save-task :add-tasklinks]
[:prepare-task :write-belief]
[:prepare-task :write-goal]
[:prepare-task :write-question]
[:prepare-task :write-quest]])
(def flow-conditions
[{:flow/from :prepare-task
:flow/to [:write-belief]
:flow/predicate ::belief?}
{:flow/from :prepare-task
:flow/to [:write-goal]
:flow/predicate ::goal?}
{:flow/from :prepare-task
:flow/to [:write-question]
:flow/predicate ::question?}
{:flow/from :prepare-task
:flow/to [:write-quest]
:flow/predicate ::quest?}])
(def kafka-io
#{:read-task :write-belief :write-goal :write-question :write-quest})
(def leaves #{:add-tasklinks})
(def build-job
(u/get-job-builder (namespace ::k) workflow kafka-io leaves flow-conditions))
;-------------------------------------------------------------------------------
;predicates
(defn goal? [_ _ {{:keys [task-type]} :message} _]
(= :goal task-type))
(defn belief? [_ _ {{:keys [task-type]} :message} _]
(= :belief task-type))
(defn question? [_ _ {{:keys [task-type]} :message} _]
(= :question task-type))
(defn quest? [_ _ {{:keys [task-type]} :message} _]
(= :quest task-type))
;-------------------------------------------------------------------------------
;functions
(defn save-task [mem task]
(info :tsk task)
(let [id (m/add-task mem task)]
(assoc task :id id)))
(defn prepare-task [_ task]
(info :prep task)
{:message task})
(defn children-terms
([statement] (children-terms statement 0))
([statement level]
(if (and (> 2 level) (coll? statement))
(let [[f & tail] statement
next-level (if (= 'conj f) level (inc level))
children (map #(children-terms % next-level) tail)]
(apply s/union (set tail) children))
[statement])))
(defn add-tasklinks [mem {:keys [statement occurrence] :as task}]
(doseq [term (conj (children-terms statement) occurrence)]
(m/add-term mem term)
(u/add-tasklink* mem term task "link")))
+52
View File
@@ -0,0 +1,52 @@
(ns narjure.control.utils
(:require [narjure.system :refer [mem inference]]
[onyx.tasks.kafka :refer [consumer producer]]))
(defn intermediate-function
[name function]
{:onyx/name name
:onyx/fn function
:onyx/type :function})
(defn leaf-function
[name function batch-settings]
(merge {:onyx/name name
:onyx/fn function
:onyx/plugin :onyx.peer.function/function
:onyx/medium :function
:onyx/type :output
:onyx/batch-size 20}
batch-settings))
(defn inject-memory
[{params :onyx.core/params} _]
{:onyx.core/params (conj params mem)})
(defn kafka-consumer
[name topic zk-address batch-settings n-peers]
(consumer
name
(merge {:kafka/topic topic
:kafka/group-id "onyx-consumer"
:kafka/zookeeper zk-address
:kafka/offset-reset :smallest
:kafka/force-reset? true
:kafka/deserializer-fn :onyx.tasks.kafka/deserialize-message-edn
:onyx/max-peers n-peers
:onyx/min-peers n-peers}
batch-settings)))
(def default-producer-config
{:kafka/serializer-fn :onyx.tasks.kafka/serialize-message-edn
:kafka/request-size 307200})
(defn kafka-producer
[name topic zk-address batch-settings]
(producer name
(merge {:kafka/topic topic
:kafka/zookeeper zk-address}
default-producer-config
batch-settings)))
(def memory-call
{:lifecycle/before-task-start inject-memory})
+11 -19
View File
@@ -2,31 +2,23 @@
(defprotocol Memory
(term [mem concept])
(select-truth [mem concept occurrence])
(truths [mem concept])
(desires [mem concept])
(beliefs [mem concept])
(goals [mem concept])
(questions [mem concept])
(quests [mem concept])
(tasklinks [mem concept])
(select-tasklink [mem concept])
(termlinks [mem concept])
(select-termlink [mem concept])
(budget [mem concept])
(select-task [mem concept])
(select-termlink [mem concept])
(select-belief [mem concept occurrence])
(add-term [mem concept])
(add-truth [mem concept truth])
(add-desire [mem concept desire])
(add-tasklink [mem concept link])
(add-task [mem task])
(add-tasklink [mem concept link type])
(add-termlink [mem concept link])
(remove-truth [mem concept id])
(remove-desire [mem concept id])
(remove-tasklink [mem concept id])
(remove-termlink [mem concept id])
(update-budget [mem update-fn])
(task [mem id])
(push-task [mem task])
(pop-task [mem])
(activate-concept [mem concept])
(pull-activated-concepts [mem]))
(increment-value [mem id name value]))
+98 -89
View File
@@ -1,42 +1,33 @@
(ns narjure.memory.redis
(:require [taoensso.carmine :as c]
[narjure.memory.api :refer [Memory]])
[narjure.memory.api :refer [Memory] :as m]
[taoensso.timbre :refer [info]])
(:import (java.util UUID)))
;postfixes for keys
(def truths-pr "_t")
(def desires-pr "_d")
(def tasklinks-pr "_tkl")
(def termlinks-pr "_tml")
(def budget-pr "_bg")
(def task-pr "_tsk")
(def tasklinks-pf "_tkl")
(def termlinks-pf "_tml")
(def budget-pf "_bg")
(def task-pf "_tsk")
(def parse-float #(Float/parseFloat %))
(def parse-int #(Integer/parseInt %))
(def parse-boolean #(Boolean/parseBoolean %))
(def truth-schema
{:frequency :float
:confidence :float
:occurrence :int
:evidences :any})
(def desire-schema
{:plausibility :float
:desirability :float
:occurrence :int
:evidences :any})
(defn parse-boolean
[val]
(info :opopop val (type val))
(Boolean/parseBoolean val))
(def task-schema
{:task-type :keyword
:evidences :vector
:eternal :boolean
:occurrence :int
:frequency :float
:confidence :float
:plausibility :float
:desirability :float
:term :any})
{:task-type :keyword
:evidences :vector
:eternal :boolean
:occurrence :int
:frequency :float
:confidence :float
:plausibility :float
:desirability :float
:statement :any
:creation-time :int})
(def termlink-schema
{:priority :float
@@ -53,7 +44,7 @@
(def deserialization-fn
{:float parse-float
:int parse-int
:boolean parse-boolean})
:keyword keyword})
(defn apply-schema [val]
(->> val
@@ -61,14 +52,12 @@
[k (get deserialization-fn v identity)]))
(into {})))
(def deserialization-map
(def deserial-map
(reduce (fn [ac [key val]] (assoc ac key (apply-schema val)))
{}
{truths-pr truth-schema
desires-pr desire-schema
task-pr task-schema
tasklinks-pr tasklink-schema
termlinks-pr termlink-schema}))
{task-pf task-schema
tasklinks-pf tasklink-schema
termlinks-pf termlink-schema}))
(defn- check-hash [val]
(if (or (integer? val) (string? val)) val (hash val)))
@@ -76,11 +65,17 @@
(defn get-key [concept postfix]
(str (check-hash concept) postfix))
(defn- get-maps-ids [conn concept postfix]
(defn- get-maps-ids-from-set [conn concept postfix]
(->> (get-key concept postfix)
c/smembers
(c/wcar conn)))
(defn- get-maps-ids-from-hash [conn concept postfix]
(->> (get-key concept postfix)
c/hgetall
(c/wcar conn)
(apply hash-map)))
(defn xf [trans-map]
(comp (partition-all 2)
(map (fn [[k v]]
@@ -91,12 +86,22 @@
(defn- get-map-by-key [conn trans-map k]
(into {} (xf trans-map) (c/wcar conn (c/hgetall k))))
(defn- get-maps [conn concept postfix]
(map (partial get-map-by-key conn (deserialization-map postfix))
(get-maps-ids conn concept postfix)))
(defn- get-maps-from-set [conn concept postfix]
(map (partial get-map-by-key conn (deserial-map postfix))
(get-maps-ids-from-set conn concept postfix)))
(defn- get-tasklinks [conn concept pred]
(let [links (->> (get-maps-ids-from-hash conn concept tasklinks-pf)
(filter (fn [[_ v]] (pred v)))
keys
(mapv (partial get-map-by-key conn (deserial-map tasklinks-pf))))]
(seq (mapv (fn [{:keys [task] :as link}]
(let [task (get-map-by-key conn (deserial-map task-pf) task)]
(assoc link :task task)))
links))))
(defn- get-map [conn concept postfix]
(get-map-by-key conn (deserialization-map postfix) (get-key concept postfix)))
(get-map-by-key conn (deserial-map postfix) (get-key concept postfix)))
(defn- add-map [conn concept postfix data]
(let [id (str (UUID/randomUUID) postfix)]
@@ -107,68 +112,72 @@
(c/wcar conn (c/srem (get-key concept postfix) id))
(c/wcar conn (c/del id)))
(defn- push [conn task]
(let [id (str (UUID/randomUUID) "_tsk")]
(c/wcar conn (c/lpush "tasks" id))
(c/wcar conn (c/hmset* id task))))
(defn remove-tasklink* [conn concept id]
(c/wcar conn (c/hdel (get-key concept tasklinks-pf) id))
(c/wcar conn (c/del id)))
(defn- tpop [conn]
(let [id (c/wcar conn (c/rpop))]
(get-map-by-key conn (deserialization-map task-pr) id)))
(defn pull-concepts [conn]
(-> (c/wcar
conn
(c/multi)
(c/smembers :active-concepts)
(println val)
(c/del :active-concepts)
(c/exec))
last
first))
(defn- activate-concept* [conn concept]
(c/wcar conn (c/sadd :active-concepts (check-hash concept))))
(defn- add-task* [conn task]
(let [id (str (UUID/randomUUID) task-pf)]
(c/wcar conn (c/hmset* id task))
id))
(defn- select-link [conn concept prefix]
(let [id (->> prefix
(get-key concept)
c/srandmember
(c/wcar conn))]
(get-map-by-key conn (deserialization-map prefix) id)))
(get-map-by-key conn (deserial-map prefix) id)))
(defn- select-task [conn concept]
(let [link (->> (c/wcar conn (c/hkeys (get-key concept tasklinks-pf)))
(mapv (partial get-map-by-key conn (deserial-map tasklinks-pf)))
rand-nth)
task-id (:task link)]
(assoc link :task (get-map-by-key conn (deserial-map task-pf) task-id))))
(defn get-task [conn id]
(get-map-by-key conn (deserialization-map task-pr) id))
(get-map-by-key conn (deserial-map task-pf) id))
(defn belief? [s] (= "belief" s))
(defn goal? [s] (= "goal" s))
(defn question? [s] (= "question" s))
(defn quest? [s] (= "quest" s))
(defn add-term* [conn concept]
(c/wcar conn (c/set (hash concept) concept)))
(defn add-tasklink* [conn concept link type]
(let [id (str (UUID/randomUUID) tasklinks-pf)]
(c/wcar conn (c/hset (get-key concept tasklinks-pf) id type))
(c/wcar conn (c/hmset* id (assoc link :id id)))
id))
(defn increment-value* [conn id name value]
(c/wcar conn (c/hincrbyfloat id name value)))
(defrecord RedisMemory
[conn]
Memory
(term [_ concept] (c/wcar conn (c/get (check-hash concept))))
(select-truth [_ concept occurrence] (select-link conn concept truths-pr))
(truths [_ concept] (get-maps conn concept truths-pr))
(desires [_ concept] (get-maps conn concept desires-pr))
(tasklinks [_ concept] (get-maps conn concept tasklinks-pr))
(select-tasklink [_ concept] (select-link conn concept tasklinks-pr))
(termlinks [_ concept] (get-maps conn concept termlinks-pr))
(select-termlink [_ concept] (select-link conn concept termlinks-pr))
(budget [_ concept] (get-map conn concept budget-pr))
(beliefs [_ concept] (get-tasklinks conn concept belief?))
(goals [_ concept] (get-tasklinks conn concept goal?))
(questions [_ concept] (get-tasklinks conn concept question?))
(quests [_ concept] (get-tasklinks conn concept quest?))
(tasklinks [_ concept] (get-tasklinks conn concept (constantly true)))
(termlinks [_ concept] (get-maps-from-set conn concept termlinks-pf))
(budget [_ concept] (get-map conn concept budget-pf))
(add-term [_ concept] (c/wcar conn (c/set (hash concept) concept)))
(add-truth [_ concept truth] (add-map conn concept truths-pr truth))
(add-desire [_ concept desire] (add-map conn concept desires-pr desire))
(add-tasklink [_ concept link] (add-map conn concept tasklinks-pr link))
(add-termlink [_ concept link] (add-map conn concept termlinks-pr link))
(select-task [_ concept] (select-task conn concept))
(select-termlink [_ concept] (select-link conn concept termlinks-pf))
(select-belief [mem concept occurrence]
(let [beliefs (m/beliefs mem concept)]
(when (sequential? beliefs) (rand-nth beliefs))))
(remove-truth [_ concept id] (remove-map conn concept truths-pr id))
(remove-desire [_ concept id] (remove-map conn concept desires-pr id))
(remove-tasklink [_ concept id] (remove-map conn concept tasklinks-pr id))
(remove-termlink [_ concept id] (remove-map conn concept termlinks-pr id))
(add-term [_ concept] (add-term* conn concept))
(add-task [_ task] (add-task* conn task))
(add-tasklink [_ concept link type] (add-tasklink* conn concept link type))
(add-termlink [_ concept link] (add-map conn concept termlinks-pf link))
(update-budget [_ update-fn])
(task [_ id] (get-task conn id))
(push-task [_ task] (push conn task))
(pop-task [_] (tpop conn))
(activate-concept [_ concept] (activate-concept* conn concept))
(pull-activated-concepts [_] (pull-concepts conn)))
(remove-tasklink [_ concept id] (remove-tasklink* conn concept id))
(remove-termlink [_ concept id] (remove-map conn concept termlinks-pf id))
(increment-value [_ id name value] (increment-value* conn id name value)))
+15 -6
View File
@@ -1,15 +1,24 @@
(ns narjure.system
(:require [mount.core :refer [defstate]]
(:require [mount.core :refer [defstate start]]
[narjure.memory.redis :as r]
[nal.deriver.rules :refer [compile-rules]]
[nal.rules :refer [all-rules]]
[nal.core :as c]))
[nal.core :as c]
[aero.core :refer [read-config]]))
(declare memory inference)
(declare mem inference)
(def config
(read-config (clojure.java.io/resource "config.edn") {:profile :dev}))
(defn get-config [path]
(get-in config path))
(def redis-config
{:spec {:host "127.0.0.1" :port 6379}})
{:spec {:uri (get-config [:redis-config :uri])}})
(defstate memory :start (r/->RedisMemory redis-config))
(defstate inference :start #(let [rules (compile-rules all-rules)]
(defstate mem :start (r/->RedisMemory redis-config))
(defstate inference :start (let [rules (compile-rules all-rules)]
(partial c/inference rules)))
(start)
-64
View File
@@ -1,64 +0,0 @@
(ns nal.test.underiver
(:require nal.reader
[nal.deriver.rules :as r]
[nal.deriver.matching :as m]
[nal.deriver.utils :as u]
[clojure.core.unify :as un]
[clojure.core.match :as omg]
[clojure.set :as cs]
[nal.core :as c]))
(r/defrules rls
#R[(P ==> M) (S ==> M) |- (S ==> P) :post (:t/induction :allow-backward) :pre ((:!= S P))])
(def compiled (r/compile-rules rls))
(def r-map (first (r/rule (first rls))))
(defn sym-map [m p1 p2]
(let [vals (set (vals m))
all (cs/difference (set (remove u/operator? (flatten [p1 p2])))
vals)]
(->> all
(map (fn [el] [`(quote ~el) el]))
(into {})
(merge m))))
(defn underiver [{:keys [p1 p2 conclusions]}]
(let [concl (vec (:conclusion (first conclusions)))
[m pattern] (m/find-and-replace-symbols concl "x")
m (sym-map m p1 p2)
p1 (m/replace-symbols p1 m)
p2 (m/replace-symbols p2 m)]
(eval (m/quote-operators
`(fn [xn#] (omg/match xn#
~pattern [~p1 ~p2]
:else []))))))
(defn deriver [rls]
(let [compiled (r/compile-rules rls)]
(fn [[p1 p2]]
(let [t {:statement p1
:desire [1 0.9]
:task-type :judgement
:occurrence 1}
b {:statement p2
:truth [1 0.9]
:occurrence 0}]
(c/inference compiled t b)))))
(comment
((underiver r-map) '[==> wut ahh?])
=> [[==> ahh? M] [==> wut M]]
((deriver rls) '[[==> ahh? M] [==> wut M]])
(let [[{st :statement}] (c/inference compiled
{:statement '[==> ahh? M]
:truth [1 0.9]
:task-type :judgement
:occurrence 1}
{:statement '[==> wut M]
:truth [1 0.9]
:occurrence 0})]
((underiver r-map) st))
)
-41
View File
@@ -1,41 +0,0 @@
(ns narjure.test.control.buffer
(:require
[clojure.test :refer :all]
[narjure.control.buffer :refer :all]
[clojure.core.async.impl.protocols :refer [full? add! remove! close-buf!]]))
(defmacro throws? [expr]
`(try
~expr
false
(catch Throwable _# true)))
(def warn-cnt (atom 0))
(defn panick! [_]
(swap! warn-cnt inc))
(deftest sliding-buffer-tests
(let [fb (panicking-sliding-buffer 2 panick! 1)]
(reset! warn-cnt 0)
(is (= 0 (count fb)))
(add! fb :1)
(is (= 1 (count fb)))
(add! fb :2)
(is (= 2 (count fb)))
(is (= 1 @warn-cnt))
(is (not (full? fb)))
(is (not (throws? (add! fb :3))))
(is (= 2 (count fb)))
(is (= :2 (remove! fb)))
(is (not (full? fb)))
(is (= 1 (count fb)))
(is (= :3 (remove! fb)))
(is (= 0 (count fb)))
(is (throws? (remove! fb)))))
-65
View File
@@ -1,65 +0,0 @@
(ns narjure.test.control.flow
(:require
[clojure.test :refer :all]
[narjure.control.flow :refer :all]
[clojure.core.async :as as]))
(deftest test-check-element-in-map
(let [m {:a 1 :b 2}]
(is (= 0 (:c (check-element-in-map [:c] m))))
(is (= [] (:c (check-element-in-map [:c] [] m))))
(is (= nil (:k (check-element-in-map [:c] [] m))))))
(deftest test-all
(is (= #{:a :b :c :d :e}
(all [[:a :b]
[:a :c]
[:c :d]
[:c :e]]))))
(defn some-var [])
(deftest test-kw->var
(is (var? (kw->fn :narjure.control.flow/all)))
(is (var? (kw->fn ::some-var)))
(is (nil? (kw->fn ::wrong-var))))
(def wf [[:a :b]
[:a :c]
[:a :d]
[:b :d]])
(deftest test-fn-outputs
(let [outputs (fn-outputs wf 1)]
(is (= 4 (count (keys outputs))))
(is (every? nil? (map outputs [:c :d])))))
(deftest test-fn-inputs
(is (= {[:d :c :b] [:a]
[:d] [:b]}
(fn-inputs wf)))
(is (= {[:d :c :b] [:a]
[:k :d] [:c :b]}
(fn-inputs (concat wf [[:c :d]
[:c :k]
[:b :k]])))))
(def out (as/chan 2))
(defn first-fn [data] (assoc data :fn1 :ok))
(defn second-fn [data] (assoc data :fn2 :ok))
(defn third-fn [data]
(as/>!! out (assoc data :fn3 :ok)))
(def test-flow [[::first-fn ::second-fn]
[::second-fn ::third-fn]
[::first-fn ::third-fn]])
(deftest test-generate-flow
(let [in (generate-flow test-flow {:buffer 2})]
(as/>!! in {})
(let [res [(as/<!! out) (as/<!! out)]]
(is (= #{{:fn1 :ok
:fn3 :ok}
{:fn1 :ok
:fn2 :ok
:fn3 :ok}}
(set res))))))
@@ -0,0 +1,88 @@
(ns narjure.test.control.general-inference
(:require [aero.core :refer [read-config]]
[clojure.test :refer [deftest is]]
[com.stuartsierra.component :as component]
[onyx api
[test-helper :refer [with-test-env]]]
[narjure.control.general-inference :as gi]
[narjure.system :refer [mem]]
[narjure.memory.api :as m]
[taoensso.carmine :refer [wcar] :as c]
[narjure.test.control.utils :refer [mock-kafka rand-id
take-values-from-topic]]))
(defn load-redis-data
[]
(wcar {} (c/flushall))
(let [concept1 '[--> P S]
concept2 '[--> K P]]
(m/add-term mem concept1)
(let [task {:task-type :goal
:occurrence 1
:plausibility 1
:desirability 0.9
:statement concept1}
task-id (m/add-task mem task)
tasklink {:priority 1
:durability 1
:quality 1
:task task-id}]
(m/add-tasklink mem concept1 tasklink "goal"))
(m/add-term mem concept2)
(m/add-termlink mem concept1 {:priority 1
:durability 1
:quality 1
:concept (hash concept2)})
(let [task {:frequency 1
:confidence 0.9
:occurrence 0
:statement concept2}
task-id (m/add-task mem task)
tasklink {:priority 1
:durability 1
:quality 1
:task task-id}]
(m/add-tasklink mem concept2 tasklink "belief"))))
(def input
(let [concept-hash (hash '[--> P S])]
[{:question-with-query-var true
:concept concept-hash}
{:concept concept-hash}]))
(def output-answer [{:some-answer true}])
(def output-tasks
[{:desire [1.0 0.8099999570846563]
:occurrence 1
:statement '[--> K S]
:task-type :goal}
{:desire [1.0 0.40499997854232817]
:occurrence 1
:statement '[--> S K]
:task-type :goal}])
(deftest general-inference-test
(let [input-topic (rand-id)
tasks-topic (rand-id)
answers-topic (rand-id)
{:keys [env-config peer-config]}
(read-config (clojure.java.io/resource "config.edn")
{:profile :test})
zk-address (get-in peer-config [:zookeeper/address])
job (gi/build-job zk-address input-topic tasks-topic answers-topic 10 1000)
mock (atom {})]
(try
(load-redis-data)
(with-test-env
[test-env [9 env-config peer-config]]
(onyx.test-helper/validate-enough-peers! test-env job)
(reset! mock (mock-kafka zk-address [[input-topic input]]))
(onyx.api/submit-job peer-config job)
(is (= (set (take-values-from-topic zk-address tasks-topic))
(set output-tasks)))
(is (= (take-values-from-topic zk-address answers-topic)
output-answer)))
(finally (swap! mock component/stop)))))
@@ -0,0 +1,156 @@
(ns narjure.test.control.local-inference.beliefs
(:require [clojure.test :refer [deftest is]]
[aero.core :refer [read-config]]
[com.stuartsierra.component :as component]
[taoensso.carmine :refer [wcar] :as c]
[narjure.system :refer [mem]]
[onyx api
[test-helper :refer [with-test-env]]]
[narjure.test.control.utils :refer [mock-kafka rand-id
take-values-from-topic]]
[narjure.control.local-inference.beliefs :refer [build-job]]
[narjure.memory.api :as m]
[narjure.test.control.utils :as u]))
(defn input []
(let [task {:statement 'cat
:occurrence 1
:frequency 0.8
:confidence 0.7
:durability 0.4
:eternal false
:creation-time 1
:evidences [1]}
id (m/add-task mem task)]
[{:task (assoc task :id id)}]))
(def output
[{:priority-diff -0.05675676676197694
:question {:durability 0.5
:id nil
:priority 0.7
:quality 1.0
:task {:eternal false
:occurrence 1
:statement 'cat
:task-type :question}}
:satisfaction 0.08108109675505559}])
(def beliefs
#{{:durability (float 0.1)
:priority (float 0.3)
:quality (float 1.0)
:task {:eternal false
:evidences [2]
:frequency (float 1.0)
:confidence (float 0.9)
:occurrence 1
:statement 'cat
:task-type :belief}}
{:durability (float 0.856)
:task {:frequency (float 0.8)
:confidence (float 0.7)
:creation-time 1
:durability "0.4"
:eternal false
:evidences [1]
:occurrence 1
:statement 'cat}}
{:durability (float 0.86767566)
:task {:confidence (float 0.9189189)
:creation-time 1
:durability "0.8676756845053482"
:eternal false
:evidences [1 2]
:frequency (float 0.9588235)
:occurrence 1
:statement 'cat}}})
(def questions
#{{:durability (float 0.5)
:priority (float 0.64324325)
:quality (float 1.0)
:task {:eternal false
:occurrence 1
:statement 'cat
:task-type :question}}})
(def goals
#{{:durability (float 0.9)
:priority (float 0.19199999)
:quality (float 1.0)
:task {:desirability (float 0.9)
:occurrence 1
:eternal false
:evidences [3]
:plausibility (float 1.0)
:statement 'cat
:task-type :goal}}})
(defn load-redis-data
[]
(wcar {} (c/flushall))
(let [concept1 'cat]
(m/add-term mem concept1)
(let [task {:task-type :belief
:occurrence 1
:frequency 1
:confidence 0.9
:statement concept1
:eternal false
:evidences [2]}
task2 {:task-type :goal
:occurrence 1
:plausibility 1
:desirability 0.9
:statement concept1
:eternal false
:evidences [3]}
task3 {:task-type :question
:occurrence 1
:statement concept1
:eternal false}
task-id (m/add-task mem task)
task2-id (m/add-task mem task2)
task3-id (m/add-task mem task3)
tasklink {:priority 0.3
:durability 0.1
:quality 1
:task task-id}
tasklink2 {:priority 0.8
:durability 0.9
:quality 1
:task task2-id}
tasklink3 {:priority 0.7
:durability 0.5
:quality 1
:task task3-id}]
(m/add-tasklink mem concept1 tasklink "belief")
(m/add-tasklink mem concept1 tasklink2 "goal")
(m/add-tasklink mem concept1 tasklink3 "question"))))
(deftest beliefs-li-inference-test
(let [tasks-topic (rand-id)
answers-topic (rand-id)
{:keys [env-config peer-config]}
(read-config (clojure.java.io/resource "config.edn")
{:profile :test})
zk-address (get-in peer-config [:zookeeper/address])
kafka-io {:input {:read-beliefs tasks-topic}
:output {:write-answers answers-topic}}
job (build-job zk-address kafka-io 10 1000)
mock (atom {})]
(try
(load-redis-data)
(with-test-env
[test-env [12 env-config peer-config]]
(onyx.test-helper/validate-enough-peers! test-env job)
(reset! mock (mock-kafka zk-address [[tasks-topic (input)]]))
(onyx.api/submit-job peer-config job)
(is (= (-> (take-values-from-topic zk-address answers-topic)
(assoc-in [0 :question :id] nil))
output))
(is (= beliefs (u/dissoc-ids (m/beliefs mem 'cat))))
(is (= questions (u/dissoc-ids (m/questions mem 'cat))))
(is (= goals (u/dissoc-ids (m/goals mem 'cat)))))
(finally (swap! mock component/stop)))))
@@ -0,0 +1,161 @@
(ns narjure.test.control.local-inference.goals
(:require [clojure.test :refer [deftest is]]
[aero.core :refer [read-config]]
[com.stuartsierra.component :as component]
[taoensso.carmine :refer [wcar] :as c]
[narjure.system :refer [mem]]
[onyx api
[test-helper :refer [with-test-env]]]
[narjure.test.control.utils :refer [mock-kafka rand-id
take-values-from-topic]]
[narjure.control.local-inference.goals :refer [build-job]]
[narjure.memory.api :as m]
[narjure.test.control.utils :as u]))
(defn input []
(let [task {:statement 'cat
:occurrence 1
:plausibility 0.85
:desirability 0.5
:durability 0.6
:priority 0.75
:creation-time 1
:evidences [1]
:eternal false}
id (m/add-task mem task)]
[{:task (assoc task :id id)}]))
(def output
[{:priority-diff -0.0636363763454526
:quest {:durability 0.5
:id nil
:priority 0.7
:quality 1.0
:task {:eternal false
:occurrence 1
:statement 'cat
:task-type :quest}}
:satisfaction 0.09090911061310525}])
(def beliefs
#{{:durability (float 0.7525)
:priority (float 0.3)
:quality (float 1.0)
:task {:confidence (float 0.9)
:frequency (float 1.0)
:eternal false
:evidences [2]
:occurrence 1
:statement 'cat
:task-type :belief}}})
(def quests
#{{:durability (float 0.5)
:priority (float 0.6363636)
:quality (float 1.0)
:task {:eternal false
:occurrence 1
:statement 'cat
:task-type :quest}}})
(def goals
#{{:durability (float 0.6)
:priority (float 0.20625)
:task {:creation-time 1
:desirability (float 0.5)
:durability "0.6"
:eternal false
:evidences [1]
:occurrence 1
:plausibility (float 0.85)
:priority "0.75"
:statement 'cat}}
{:durability (float 0.6363636)
:priority (float 0.20625)
:task {:creation-time 1
:desirability (float 0.9090909)
:durability "0.636363644245242"
:eternal false
:evidences [1 3]
:occurrence 1
:plausibility (float 0.26500002)
:priority "0.20624999105930325"
:statement 'cat}}
{:durability (float 0.9)
:priority (float 0.8)
:quality (float 1.0)
:task {:desirability (float 0.9)
:eternal false
:evidences [3]
:occurrence 1
:plausibility (float 0.2)
:statement 'cat
:task-type :goal}}})
(defn load-redis-data
[]
(wcar {} (c/flushall))
(let [concept1 'cat]
(m/add-term mem concept1)
(let [task {:task-type :belief
:occurrence 1
:frequency 1
:confidence 0.9
:statement concept1
:eternal false
:evidences [2]}
task2 {:task-type :goal
:occurrence 1
:plausibility 0.2
:desirability 0.9
:statement concept1
:eternal false
:evidences [3]}
task3 {:task-type :quest
:occurrence 1
:statement concept1
:eternal false}
task-id (m/add-task mem task)
task2-id (m/add-task mem task2)
task3-id (m/add-task mem task3)
tasklink {:priority 0.3
:durability 0.1
:quality 1
:task task-id}
tasklink2 {:priority 0.8
:durability 0.9
:quality 1
:task task2-id}
tasklink3 {:priority 0.7
:durability 0.5
:quality 1
:task task3-id}]
(m/add-tasklink mem concept1 tasklink "belief")
(m/add-tasklink mem concept1 tasklink2 "goal")
(m/add-tasklink mem concept1 tasklink3 "quest"))))
(deftest goals-li-inference-test
(let [tasks-topic (rand-id)
answers-topic (rand-id)
{:keys [env-config peer-config]}
(read-config (clojure.java.io/resource "config.edn")
{:profile :test})
zk-address (get-in peer-config [:zookeeper/address])
kafka-io {:input {:read-goals tasks-topic}
:output {:write-answers answers-topic}}
job (build-job zk-address kafka-io 10 1000)
mock (atom {})]
(try
(load-redis-data)
(with-test-env
[test-env [12 env-config peer-config]]
(onyx.test-helper/validate-enough-peers! test-env job)
(reset! mock (mock-kafka zk-address [[tasks-topic (input)]]))
(onyx.api/submit-job peer-config job)
(is (= (-> (take-values-from-topic zk-address answers-topic)
(assoc-in [0 :quest :id] nil))
output))
(is (= beliefs (u/dissoc-ids (m/beliefs mem 'cat))))
(is (= quests (u/dissoc-ids (m/quests mem 'cat))))
(is (= goals (u/dissoc-ids (m/goals mem 'cat)))))
(finally (swap! mock component/stop)))))
@@ -0,0 +1,114 @@
(ns narjure.test.control.local-inference.questions
(:require [clojure.test :refer [deftest is]]
[aero.core :refer [read-config]]
[com.stuartsierra.component :as component]
[taoensso.carmine :refer [wcar] :as c]
[narjure.system :refer [mem]]
[onyx api
[test-helper :refer [with-test-env]]]
[narjure.test.control.utils :refer [mock-kafka rand-id
take-values-from-topic]]
[narjure.control.local-inference.questions :refer [build-job]]
[narjure.memory.api :as m]
[narjure.test.control.utils :as u]))
(def input
[{:task {:statement 'cat
:occurrence 1
:durability 1
:priority 1
:evidences [1]
:eternal false}}])
(def output
[{:belief {:durability 0.78
:id nil
:priority 0.56
:quality 1.0
:task {:confidence 0.9
:frequency 0.8
:eternal false
:evidences [2]
:occurrence 1
:statement 'cat
:task-type :belief}}
:durability-diff 0.02200000810623237
:question {:durability 1
:occurrence 1
:priority 0.8999999761581421
:eternal false
:evidences [1]
:statement 'cat}}])
(def beliefs
#{{:durability (float 0.802)
:priority (float 0.56)
:quality (float 1.0)
:task {:confidence (float 0.9)
:frequency (float 0.8)
:occurrence 1
:eternal false
:evidences [2]
:statement 'cat
:task-type :belief}}})
(def questions
#{{:durability (float 1.0)
:priority (float 0.9)
:task {}}})
(defn load-redis-data
[]
(wcar {} (c/flushall))
(let [concept1 'cat]
(m/add-term mem concept1)
(let [task {:task-type :belief
:occurrence 1
:frequency 0.8
:confidence 0.9
:statement concept1
:evidences [2]
:eternal false}
task2 {:task-type :goal
:occurrence 1
:plausibility 1
:desirability 0.9
:statement concept1
:evidences [3]
:eternal false}
task-id (m/add-task mem task)
task2-id (m/add-task mem task2)
tasklink {:priority 0.56
:durability 0.78
:quality 1
:task task-id}
tasklink2 {:priority 0.65
:durability 0.8
:quality 1
:task task2-id}]
(m/add-tasklink mem concept1 tasklink "belief")
(m/add-tasklink mem concept1 tasklink2 "goal"))))
(deftest questions-li-inference-test
(let [tasks-topic (rand-id)
answers-topic (rand-id)
{:keys [env-config peer-config]}
(read-config (clojure.java.io/resource "config.edn")
{:profile :test})
zk-address (get-in peer-config [:zookeeper/address])
kafka-io {:input {:read-questions tasks-topic}
:output {:write-answer answers-topic}}
job (build-job zk-address kafka-io 10 1000)
mock (atom {})]
(try
(load-redis-data)
(with-test-env
[test-env [6 env-config peer-config]]
(onyx.test-helper/validate-enough-peers! test-env job)
(reset! mock (mock-kafka zk-address [[tasks-topic input]]))
(onyx.api/submit-job peer-config job)
(is (= (-> (take-values-from-topic zk-address answers-topic)
(assoc-in [0 :belief :id] nil))
output))
(is (= beliefs (u/dissoc-ids (m/beliefs mem 'cat))))
(is (= questions (u/dissoc-ids (m/questions mem 'cat)))))
(finally (swap! mock component/stop)))))
@@ -0,0 +1,123 @@
(ns narjure.test.control.local-inference.quests
(:require [clojure.test :refer [deftest is]]
[aero.core :refer [read-config]]
[com.stuartsierra.component :as component]
[taoensso.carmine :refer [wcar] :as c]
[narjure.system :refer [mem]]
[onyx api
[test-helper :refer [with-test-env]]]
[narjure.test.control.utils :refer [mock-kafka rand-id
take-values-from-topic]]
[narjure.control.local-inference.quests :refer [build-job]]
[narjure.memory.api :as m]
[narjure.test.control.utils :as u]))
(defn input []
(let [task {:statement 'cat
:occurrence 1
:durability 0.95
:priority 1
:evidences [1]
:eternal false}
id (m/add-task mem task)]
[{:task (assoc task :id id)}]))
(def output
[{:durability-diff 0.020000003576278402
:goal {:durability 0.8
:id nil
:priority 0.65
:quality 1.0
:task {:desirability 0.9
:occurrence 1
:eternal false
:evidences [3]
:plausibility 1.0
:statement 'cat
:task-type :goal}}
:quest {:id nil
:durability 0.95
:occurrence 1
:priority 0.8999999761581421
:eternal false
:evidences [1]
:statement 'cat}}])
(def goals
#{{:durability (float 0.82)
:priority (float 0.65)
:quality (float 1.0)
:task {:desirability (float 0.9)
:occurrence 1
:eternal false
:evidences [3]
:plausibility (float 1.0)
:statement 'cat
:task-type :goal}}})
(def quests
#{{:durability (float 0.95)
:priority (float 0.9)
:task {:durability "0.95"
:eternal false
:evidences [1]
:occurrence 1
:priority "1"
:statement 'cat}}})
(defn load-redis-data
[]
(wcar {} (c/flushall))
(let [concept1 'cat]
(m/add-term mem concept1)
(let [task {:task-type :belief
:occurrence 1
:frequency 0.8
:confidence 0.9
:statement concept1
:evidences [2]
:eternal false}
task2 {:task-type :goal
:occurrence 1
:plausibility 1
:desirability 0.9
:statement concept1
:evidences [3]
:eternal false}
task-id (m/add-task mem task)
task2-id (m/add-task mem task2)
tasklink {:priority 0.56
:durability 0.78
:quality 1
:task task-id}
tasklink2 {:priority 0.65
:durability 0.8
:quality 1
:task task2-id}]
(m/add-tasklink mem concept1 tasklink "belief")
(m/add-tasklink mem concept1 tasklink2 "goal"))))
(deftest quests-li-inference-test
(let [tasks-topic (rand-id)
answers-topic (rand-id)
{:keys [env-config peer-config]}
(read-config (clojure.java.io/resource "config.edn")
{:profile :test})
zk-address (get-in peer-config [:zookeeper/address])
kafka-io {:input {:read-quests tasks-topic}
:output {:write-answer answers-topic}}
job (build-job zk-address kafka-io 10 1000)
mock (atom {})]
(try
(load-redis-data)
(with-test-env
[test-env [6 env-config peer-config]]
(onyx.test-helper/validate-enough-peers! test-env job)
(reset! mock (mock-kafka zk-address [[tasks-topic (input)]]))
(onyx.api/submit-job peer-config job)
(is (= (-> (take-values-from-topic zk-address answers-topic)
(assoc-in [0 :goal :id] nil)
(assoc-in [0 :quest :id] nil))
output))
(is (= goals (u/dissoc-ids (m/goals mem 'cat))))
(is (= quests (u/dissoc-ids (m/quests mem 'cat)))))
(finally (swap! mock component/stop)))))
@@ -0,0 +1,91 @@
(ns narjure.test.control.task-preprocessing
(:require
[clojure.test :refer [deftest is]]
[aero.core :refer [read-config]]
[com.stuartsierra.component :as component]
[taoensso.carmine :refer [wcar] :as c]
[narjure.system :refer [mem]]
[onyx api
[test-helper :refer [with-test-env]]]
[narjure.test.control.utils :refer [mock-kafka rand-id
take-values-from-topic]]
[narjure.control.task-preprocessing :refer [build-job]]
[narjure.memory.api :as m]
[narjure.test.control.utils :as u]))
(def input
[{:statement '[--> cat aimal]
:task-type :belief
:occurrence 1}
{:statement '[--> tim cat]
:task-type :question
:occurrence 1}
{:statement '[--> computer intelligent]
:task-type :goal
:occurrence 1}
{:statement '[<=> wut duck]
:task-type :quest
:occurrence 1}])
(def beliefs
#{{:occurrence 1
:statement '[--> cat aimal]
:task-type :belief}})
(def questions
#{{:statement '[--> tim cat]
:task-type :question
:occurrence 1}})
(def goals
#{{:statement '[--> computer intelligent]
:task-type :goal
:occurrence 1}})
(def quests
#{{:statement '[<=> wut duck]
:task-type :quest
:occurrence 1}})
(def cat-tasklinks
#{{:task {:occurrence 1
:statement '[--> cat aimal]
:task-type :belief}}
{:task {:occurrence 1
:statement '[--> tim cat]
:task-type :question}}})
(deftest quests-li-inference-test
(let [read-tasks (rand-id)
write-beliefs (rand-id)
write-goals (rand-id)
write-questions (rand-id)
write-quests (rand-id)
{:keys [env-config peer-config]}
(read-config (clojure.java.io/resource "config.edn")
{:profile :test})
zk-address (get-in peer-config [:zookeeper/address])
kafka-io {:input {:read-task read-tasks}
:output {:write-belief write-beliefs
:write-goal write-goals
:write-question write-questions
:write-quest write-quests}}
job (build-job zk-address kafka-io 10 1000)
mock (atom {})]
(try
(with-test-env
[test-env [8 env-config peer-config]]
(onyx.test-helper/validate-enough-peers! test-env job)
(reset! mock (mock-kafka zk-address [[read-tasks input]]))
(onyx.api/submit-job peer-config job)
(wcar {} (c/flushall))
(is (= (u/dissoc-ids (take-values-from-topic zk-address write-beliefs))
beliefs))
(is (= (u/dissoc-ids (take-values-from-topic zk-address write-goals))
goals))
(is (= (u/dissoc-ids (take-values-from-topic zk-address write-questions))
questions))
(is (= (u/dissoc-ids (take-values-from-topic zk-address write-quests))
quests))
(is (= (u/dissoc-ids (m/tasklinks mem 'cat)) cat-tasklinks)))
(finally (swap! mock component/stop)))))
+42
View File
@@ -0,0 +1,42 @@
(ns narjure.test.control.utils
(:require [onyx.kafka.embedded-server :as ke]
[com.stuartsierra.component :as component]
[clj-kafka
[admin :as kadmin]
[producer :as kp]]
[onyx.kafka.utils :refer [take-until-done]])
(:import (java.util UUID)))
(defn rand-id [] (str "onyx-test-" (UUID/randomUUID)))
(defn mock-kafka
[zookeeper inputs]
(let [log-dir (str "/tmp/embedded-kafka" (UUID/randomUUID))
kafka-server (component/start
(ke/map->EmbeddedKafka
{:hostname "127.0.0.1"
:port 9092
:broker-id 0
:log-dir log-dir
:zookeeper-addr zookeeper}))
producer1 (kp/producer
{"metadata.broker.list" "127.0.0.1:9092"
"serializer.class" "kafka.serializer.DefaultEncoder"
"partitioner.class" "kafka.producer.DefaultPartitioner"})]
(doseq [[topic input] inputs]
(do (doseq [x (concat input [:done])]
(->> (pr-str x)
.getBytes
(kp/message topic)
(kp/send-message producer1)))))
kafka-server))
(defn take-values-from-topic [zk-address topic]
(->> (take-until-done zk-address topic #(read-string (String. % "UTF-8")))
(sort-by (comp :n :value))
(mapv :value)))
(defn dissoc-ids
[coll]
(set (mapv #(dissoc % :id) coll)))
+30 -13
View File
@@ -13,13 +13,20 @@
(def concept1 '[--> tim cat])
(def concept2 'tim)
(def truth1 {:frequency (float 0.9)
:confidence (float 0.2)
:occurrence 1})
(def belief1 {:frequency (float 0.9)
:confidence (float 0.2)
:occurrence 1
:term concept1})
(def truth2 {:frequency (float 0.5)
:confidence (float 0.8)
:occurrence 2})
(def belief2 {:frequency (float 0.5)
:confidence (float 0.8)
:occurrence 2
:term concept1})
(def goal1 {:frequency (float 0.5)
:confidence (float 0.8)
:occurrence 2
:term concept1})
(defn termlink [c]
{:priority (float 0.1)
@@ -27,21 +34,31 @@
:quality (float 0.1)
:concept (str (hash c))})
(defn tasklink [task-id]
{:priority (float 0.1)
:durability (float 0.1)
:quality (float 0.1)
:task task-id})
(deftest test-redis
(c/wcar config (c/flushall))
(m/add-term mem concept1)
(is (= concept1 (m/term mem (hash concept1))))
(m/add-truth mem concept1 truth1)
(m/add-truth mem concept1 truth2)
(let [id1 (m/add-task mem belief1)
id2 (m/add-task mem belief2)
id3 (m/add-task mem goal1)]
(m/add-tasklink mem concept1 (tasklink id1) "belief")
(m/add-tasklink mem concept1 (tasklink id2) "belief")
(m/add-tasklink mem concept1 (tasklink id3) "goal")
(is (= (set [belief1 belief2])
(set (map :task (m/beliefs mem concept1)))))
(is (= (set [truth1 truth2])
(set (map #(dissoc % :id) (m/truths mem concept1)))))
(m/remove-truth mem concept1 (:id (last (m/truths mem concept1))))
(is (= [goal1]
(map :task (m/goals mem concept1)))))
(is (:id (first (m/truths mem concept1))))
(is ((set [truth1 truth2]) (dissoc (first (m/truths mem concept1)) :id)))
(is (:task (first (m/beliefs mem concept1))))
(m/add-term mem concept2)
(m/add-termlink mem concept1 (termlink concept2))