Compare commits
5
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d48a3379fc | ||
|
|
76cb5536c1 | ||
|
|
faf010a846 | ||
|
|
6508e881bb | ||
|
|
91bcaa0d0b |
@@ -15,4 +15,3 @@ pom.xml.asc
|
||||
\#*\#
|
||||
.\#*
|
||||
/src/nal/experiments.clj
|
||||
*.log
|
||||
|
||||
-14
@@ -1,14 +0,0 @@
|
||||
language: clojure
|
||||
|
||||
# Notify #nars
|
||||
notifications:
|
||||
slack:
|
||||
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
|
||||
+2
-11
@@ -9,14 +9,7 @@
|
||||
[org.clojure/tools.nrepl "0.2.12"]
|
||||
[org.clojure/data.priority-map "0.0.7"]
|
||||
[org.clojure/core.match "0.3.0-alpha4"]
|
||||
[org.clojure/core.unify "0.5.5"]
|
||||
[org.onyxplatform/onyx "0.9.0"]
|
||||
[org.onyxplatform/onyx-kafka "0.9.0.1"]
|
||||
[com.taoensso/carmine "2.12.2"]
|
||||
[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"]]
|
||||
[org.clojure/core.unify "0.5.5"]]
|
||||
:main ^:skip-aot narjure.core
|
||||
:plugins [[lein-cloverage "1.0.6"]
|
||||
[jonase/eastwood "0.2.3"]
|
||||
@@ -26,6 +19,4 @@
|
||||
:target-path "target/%s"
|
||||
:repl-options {:init-ns narjure.repl
|
||||
:nrepl-middleware [narjure.repl/narsese-handler]}
|
||||
:profiles {:uberjar {:aot :all}
|
||||
:test {:dependencies [[com.stuartsierra/component "0.3.1"]]}
|
||||
:dev {:dependencies [[com.stuartsierra/component "0.3.1"]]}})
|
||||
:profiles {:uberjar {:aot :all}})
|
||||
|
||||
@@ -1,29 +0,0 @@
|
||||
{: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"}}
|
||||
+9
-3
@@ -1,12 +1,18 @@
|
||||
(ns nal.core
|
||||
(:require [nal.deriver.truth :as t]
|
||||
[nal.deriver :refer [generate-conclusions]]))
|
||||
[nal.deriver :refer [generate-conclusions]]
|
||||
[nal.rules :as r]))
|
||||
|
||||
(defn choice [[f1 c1] [f2 c2]]
|
||||
(if (>= c1 c2) [f1 c1] [f2 c2]))
|
||||
|
||||
(defn inference
|
||||
[rules {:keys [task-type] :as task} belief]
|
||||
(generate-conclusions (rules task-type) task belief))
|
||||
[{:keys [task-type] :as task} belief]
|
||||
(generate-conclusions (r/rules task-type) task belief))
|
||||
|
||||
(def revision t/revision)
|
||||
|
||||
(comment
|
||||
:shift-occurrence-forward ;pre
|
||||
:shift-occurrence-backward ;pre
|
||||
:linkage-temporal)
|
||||
|
||||
+1
-2
@@ -8,8 +8,7 @@
|
||||
(defn get-matcher [rules p1 p2]
|
||||
(let [matchers (->> (mall-paths p1 p2)
|
||||
(filter rules)
|
||||
(map rules)
|
||||
(map (fn [el] (:matcher el))))]
|
||||
(map rules))]
|
||||
(case (count matchers)
|
||||
0 (constantly [])
|
||||
1 (first matchers)
|
||||
|
||||
@@ -19,7 +19,7 @@
|
||||
#{`= `not= `seq? `first `and `let `pos? `> `>= `< `<= `coll? `set `quote
|
||||
`count 'aops `- `not-empty-diff? `not-empty-inter? `walk `munification-map
|
||||
`substitute `sets `some `deref `do `vreset! `volatile! `fn `mapv `if
|
||||
`sort-commutative `n/reduce-ext-inter `n/reduce-symilarity `complement
|
||||
`sort-commutative `n/reduce-ext-inter `n/reduce-similarity `complement
|
||||
`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
|
||||
@@ -67,19 +67,24 @@
|
||||
(defn traverse-node
|
||||
"Generates code for precondition node."
|
||||
[vars result {:keys [conclusions children condition]}]
|
||||
`(when ~(quote-operators condition)
|
||||
~(when-not (zero? (count conclusions))
|
||||
`(vswap! ~result concat
|
||||
~@(set (map #(mapv (partial form-conclusion vars) %)
|
||||
(quote-operators conclusions)))))
|
||||
~@(map (fn [n] (traverse-node vars result n)) children)))
|
||||
(let [conclusions (remove
|
||||
nil?
|
||||
[(when-not (zero? (count conclusions))
|
||||
`(vswap! ~result concat
|
||||
~@(set (map #(mapv (partial form-conclusion vars) %)
|
||||
(quote-operators conclusions)))))])
|
||||
children (mapcat (fn [n] (traverse-node vars result n)) children)]
|
||||
(if (true? condition)
|
||||
(concat conclusions children)
|
||||
[`(when ~(quote-operators condition)
|
||||
~@(concat conclusions children))])))
|
||||
|
||||
(defn traversal
|
||||
"Walk through preconditions tree and generates code for matcher."
|
||||
[vars tree]
|
||||
(let [results (gensym)]
|
||||
`(let [~results (volatile! [])]
|
||||
~(traverse-node vars results tree)
|
||||
~@(traverse-node vars results tree)
|
||||
@~results)))
|
||||
|
||||
(defn replace-occurrences
|
||||
@@ -340,6 +345,5 @@
|
||||
match-fn-code (-> main-pattern
|
||||
(gen-rules rules)
|
||||
(match-rules main-pattern task-type))]
|
||||
[k (assoc v :matcher (eval match-fn-code)
|
||||
:matcher-code match-fn-code)])))
|
||||
[k (eval match-fn-code)])))
|
||||
(into {})))
|
||||
|
||||
@@ -111,7 +111,7 @@
|
||||
[_ ['ext-set & l1] ['ext-set & l2]] (diff 'ext-set l1 l2)
|
||||
:else st))
|
||||
|
||||
(defn reduce-symilarity
|
||||
(defn reduce-similarity
|
||||
[st]
|
||||
(m/match st
|
||||
['<-> ['ext-set s] ['ext-set p]] ['<-> s p]
|
||||
@@ -169,7 +169,7 @@
|
||||
'| `reduce-int-inter
|
||||
'- `reduce-ext-dif
|
||||
'int-dif `reduce-int-dif
|
||||
'<-> `reduce-symilarity
|
||||
'<-> `reduce-similarity
|
||||
'* `reduce-production
|
||||
'int-image `reduce-image
|
||||
'ext-image `reduce-image
|
||||
@@ -185,7 +185,7 @@
|
||||
| (reduce-int-inter st)
|
||||
- (reduce-ext-dif st)
|
||||
int-dif (reduce-int-dif st)
|
||||
<-> (reduce-symilarity st)
|
||||
<-> (reduce-similarity st)
|
||||
* (reduce-production st)
|
||||
int-image (reduce-image st)
|
||||
ext-image (reduce-image st)
|
||||
|
||||
@@ -5,7 +5,7 @@
|
||||
[nal.deriver.substitution :refer [substitute munification-map]]
|
||||
[nal.deriver.terms-permutation :refer [implications equivalences]]
|
||||
[clojure.set :refer [union intersection]]
|
||||
[narjure.defaults :refer [temporal-window-duration]]
|
||||
[narjure.defaults :refer [duration]]
|
||||
[clojure.core.match :as m]
|
||||
[nal.deriver.normalization :refer [reduce-seq-conj]]))
|
||||
|
||||
@@ -88,11 +88,11 @@
|
||||
[_]
|
||||
[`(not= :eternal :t-occurrence)
|
||||
`(not= :eternal :b-occurrence)
|
||||
`(<= ~temporal-window-duration (abs (- :t-occurrence :b-occurrence)))])
|
||||
`(<= ~duration (abs (- :t-occurrence :b-occurrence)))])
|
||||
|
||||
(defmethod compound-precondition :concurrent
|
||||
[_]
|
||||
[`(> ~temporal-window-duration (abs (- :t-occurrence :b-occurrence)))])
|
||||
[`(> ~duration (abs (- :t-occurrence :b-occurrence)))])
|
||||
|
||||
;-------------------------------------------------------------------------------
|
||||
(defmulti precondition-transformation (fn [arg1 _] (first arg1)))
|
||||
@@ -164,11 +164,11 @@
|
||||
(m/match (mapv #(if (and (coll? %) (= 'quote (first %)))
|
||||
(second %) %) (rest args))
|
||||
[(:or '=|> '==>)] concl
|
||||
['pred-impl] `(let [:t-occurrence (+ :t-occurrence ~temporal-window-duration)] ~concl)
|
||||
['retro-impl] `(let [:t-occurrence (- :t-occurrence ~temporal-window-duration)] ~concl)
|
||||
['pred-impl] `(let [:t-occurrence (+ :t-occurrence ~duration)] ~concl)
|
||||
['retro-impl] `(let [:t-occurrence (- :t-occurrence ~duration)] ~concl)
|
||||
[sym (:or '=|> '==>)] (shift-forward-let sym concl)
|
||||
[sym 'pred-impl] (shift-forward-let sym `+ concl temporal-window-duration)
|
||||
[sym 'retro-impl] (shift-forward-let sym `- concl temporal-window-duration)))
|
||||
[sym 'pred-impl] (shift-forward-let sym `+ concl duration)
|
||||
[sym 'retro-impl] (shift-forward-let sym `- concl duration)))
|
||||
|
||||
(defn backward-interval-check [sym]
|
||||
`(and (coll? ~sym) (= (first ~sym) (quote ~'seq-conj))
|
||||
@@ -190,7 +190,7 @@
|
||||
|
||||
(defmethod conclusion-transformation :shift-occurrence-backward
|
||||
[args concl]
|
||||
(let [duration (- temporal-window-duration)]
|
||||
(let [duration (- duration)]
|
||||
(m/match (mapv #(if (and (coll? %) (= 'quote (first %)))
|
||||
(second %) %) (rest args))
|
||||
[(:or '=|> '==>)] concl
|
||||
|
||||
@@ -57,6 +57,12 @@
|
||||
[{:keys [pre]}]
|
||||
(some #{:question?} pre))
|
||||
|
||||
(defn quest?
|
||||
"Return true if rule allows only quest as task."
|
||||
[{:keys [pre] [{post :post}] :conclusions}]
|
||||
(and (some #{:question?} pre)
|
||||
(every? #(not (#{:p/judgement} %)) post)))
|
||||
|
||||
(defn goal?
|
||||
"Return true if rule allows only goal as task."
|
||||
[{pre :pre [{post :post}] :conclusions}]
|
||||
@@ -134,10 +140,13 @@
|
||||
allow-backward? expand-backward-rules)
|
||||
judgement-rules# (check-duplication (filter judgement? rules))
|
||||
question-rules# (check-duplication (filter question? rules))
|
||||
goal-rules# (check-duplication (filter goal? rules))]
|
||||
(println "Q rules:" (count question-rules#))
|
||||
(println "J rules:" (count judgement-rules#))
|
||||
(println "G rules:" (count goal-rules#))
|
||||
goal-rules# (check-duplication (filter goal? rules))
|
||||
quest-rules# (check-duplication (filter quest? rules))]
|
||||
(println "Beliefs rules:" (count judgement-rules#))
|
||||
(println "Questions rules:" (count question-rules#))
|
||||
(println "Goal rules:" (count goal-rules#))
|
||||
(println "Quests rules:" (count quest-rules#))
|
||||
{:judgement (rules-map judgement-rules# :judgement)
|
||||
:question (rules-map question-rules# :question)
|
||||
:goal (rules-map goal-rules# :goal)})))
|
||||
:goal (rules-map goal-rules# :goal)
|
||||
:quest (rules-map quest-rules# :quest)})))
|
||||
|
||||
@@ -3,7 +3,6 @@
|
||||
|
||||
;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))
|
||||
@@ -109,6 +108,8 @@
|
||||
|
||||
(defn t-identity [p1 _] p1)
|
||||
|
||||
(defn d-identity [p1 _] p1)
|
||||
|
||||
(defn belief-identity [p1 p2] (when p2 p1))
|
||||
|
||||
(defn belief-structural-deduction [_ p2]
|
||||
@@ -132,9 +133,6 @@
|
||||
[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
|
||||
@@ -170,6 +168,6 @@
|
||||
:d/deduction intersection
|
||||
:d/weak desire-weak
|
||||
:d/induction desire-induction
|
||||
:d/identity identity
|
||||
:d/identity d-identity
|
||||
:d/negation negation
|
||||
:d/structural-strong desire-structural-strong})
|
||||
|
||||
@@ -450,3 +450,18 @@
|
||||
; compound composition one premise
|
||||
#R[(|| B :list/A) B |- (|| B :list/A) :pre (:question?) :post (:t/belief-structural-deduction :p/judgement)]
|
||||
)
|
||||
|
||||
(def rules (compile-rules all-rules))
|
||||
|
||||
(defn freq [task-type]
|
||||
"Check frequency"
|
||||
(into {} (map (fn [[k v]] [(str k) (count (:rules v))]) (task-type rules))))
|
||||
|
||||
(defn stats [task-type]
|
||||
(let [fr (freq task-type)]
|
||||
(println "Total" (reduce + (vals fr)))
|
||||
(println "Total keys" (count (task-type rules)))
|
||||
(println "Freq" (sort (frequencies (vals fr))))
|
||||
(println "Min" (reduce min (vals fr)))
|
||||
(println "Max" (reduce (fn [[_ v1 :as p] [_ v :as n]]
|
||||
(if (> v1 v) p n)) fr))))
|
||||
|
||||
@@ -1,128 +0,0 @@
|
||||
(ns narjure.control.general-inference
|
||||
(: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]]))
|
||||
|
||||
(def workflow
|
||||
[[:select-active-concepts :select-task-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 :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 question-with-query-var? [_ _ segment _]
|
||||
(if (:question-with-query-var segment) true false))
|
||||
|
||||
(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)]))
|
||||
|
||||
(def tasks-with-redis-access
|
||||
(disj (set (flatten workflow))
|
||||
:select-active-concepts
|
||||
:write-tasks
|
||||
:write-answer
|
||||
:do-general-inference))
|
||||
|
||||
(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)
|
||||
belief (m/select-belief mem linked-concept occurrence)]
|
||||
(info :belief belief)
|
||||
(assoc segment :belief belief)))
|
||||
|
||||
(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]
|
||||
:truth [frequency confidence]
|
||||
:task-type task-type
|
||||
:occurrence occurrence}
|
||||
|
||||
{:keys [confidence frequency]} belief
|
||||
belief (assoc belief :truth [frequency confidence])]
|
||||
(info :inference task belief)
|
||||
(mapv (fn [conclusion] {:message conclusion})
|
||||
(inference task belief))))
|
||||
|
||||
(defn choose-answer [mem segment]
|
||||
(info :choose-answer)
|
||||
{:message {:some-answer true}})
|
||||
@@ -1,129 +0,0 @@
|
||||
(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])
|
||||
@@ -1,135 +0,0 @@
|
||||
(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])
|
||||
@@ -1,57 +0,0 @@
|
||||
(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))
|
||||
@@ -1,57 +0,0 @@
|
||||
(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))
|
||||
@@ -1,197 +0,0 @@
|
||||
(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)))
|
||||
|
||||
@@ -1,79 +0,0 @@
|
||||
(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")))
|
||||
@@ -1,52 +0,0 @@
|
||||
(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})
|
||||
@@ -32,4 +32,4 @@
|
||||
|
||||
(def ^{:type double} horizon 1)
|
||||
|
||||
(def temporal-window-duration 80)
|
||||
(def duration 80)
|
||||
|
||||
@@ -1,24 +0,0 @@
|
||||
(ns narjure.memory.api)
|
||||
|
||||
(defprotocol Memory
|
||||
(term [mem concept])
|
||||
(beliefs [mem concept])
|
||||
(goals [mem concept])
|
||||
(questions [mem concept])
|
||||
(quests [mem concept])
|
||||
(tasklinks [mem concept])
|
||||
(termlinks [mem concept])
|
||||
(budget [mem concept])
|
||||
|
||||
(select-task [mem concept])
|
||||
(select-termlink [mem concept])
|
||||
(select-belief [mem concept occurrence])
|
||||
|
||||
(add-term [mem concept])
|
||||
(add-task [mem task])
|
||||
(add-tasklink [mem concept link type])
|
||||
(add-termlink [mem concept link])
|
||||
|
||||
(remove-tasklink [mem concept id])
|
||||
(remove-termlink [mem concept id])
|
||||
(increment-value [mem id name value]))
|
||||
@@ -1,183 +0,0 @@
|
||||
(ns narjure.memory.redis
|
||||
(:require [taoensso.carmine :as c]
|
||||
[narjure.memory.api :refer [Memory] :as m]
|
||||
[taoensso.timbre :refer [info]])
|
||||
(:import (java.util UUID)))
|
||||
|
||||
;postfixes for keys
|
||||
(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 %))
|
||||
(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
|
||||
:statement :any
|
||||
:creation-time :int})
|
||||
|
||||
(def termlink-schema
|
||||
{:priority :float
|
||||
:durability :float
|
||||
:quality :float
|
||||
:concept :string})
|
||||
|
||||
(def tasklink-schema
|
||||
{:priority :float
|
||||
:durability :float
|
||||
:quality :float
|
||||
:task :string})
|
||||
|
||||
(def deserialization-fn
|
||||
{:float parse-float
|
||||
:int parse-int
|
||||
:keyword keyword})
|
||||
|
||||
(defn apply-schema [val]
|
||||
(->> val
|
||||
(map (fn [[k v]]
|
||||
[k (get deserialization-fn v identity)]))
|
||||
(into {})))
|
||||
|
||||
(def deserial-map
|
||||
(reduce (fn [ac [key val]] (assoc ac key (apply-schema val)))
|
||||
{}
|
||||
{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)))
|
||||
|
||||
(defn get-key [concept postfix]
|
||||
(str (check-hash 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]]
|
||||
(let [k (keyword k)
|
||||
tf (get trans-map k identity)]
|
||||
[k (tf v)])))))
|
||||
|
||||
(defn- get-map-by-key [conn trans-map k]
|
||||
(into {} (xf trans-map) (c/wcar conn (c/hgetall k))))
|
||||
|
||||
(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 (deserial-map postfix) (get-key concept postfix)))
|
||||
|
||||
(defn- add-map [conn concept postfix data]
|
||||
(let [id (str (UUID/randomUUID) postfix)]
|
||||
(c/wcar conn (c/sadd (get-key concept postfix) id))
|
||||
(c/wcar conn (c/hmset* id (assoc data :id id)))))
|
||||
|
||||
(defn- remove-map [conn concept postfix id]
|
||||
(c/wcar conn (c/srem (get-key concept postfix) id))
|
||||
(c/wcar conn (c/del id)))
|
||||
|
||||
(defn remove-tasklink* [conn concept id]
|
||||
(c/wcar conn (c/hdel (get-key concept tasklinks-pf) id))
|
||||
(c/wcar conn (c/del id)))
|
||||
|
||||
(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 (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 (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))))
|
||||
(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))
|
||||
|
||||
(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))))
|
||||
|
||||
(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))
|
||||
|
||||
(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)))
|
||||
@@ -1,24 +0,0 @@
|
||||
(ns narjure.system
|
||||
(: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]
|
||||
[aero.core :refer [read-config]]))
|
||||
|
||||
(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 {:uri (get-config [:redis-config :uri])}})
|
||||
|
||||
(defstate mem :start (r/->RedisMemory redis-config))
|
||||
(defstate inference :start (let [rules (compile-rules all-rules)]
|
||||
(partial c/inference rules)))
|
||||
|
||||
(start)
|
||||
+72
-185
@@ -1,28 +1,25 @@
|
||||
(ns nal.test.core
|
||||
(:require [clojure.test :refer :all]
|
||||
[nal.core :as c]
|
||||
[nal.rules :as r]
|
||||
[nal.deriver.rules :refer [compile-rules]]))
|
||||
|
||||
(def inference (partial c/inference (compile-rules r/all-rules)))
|
||||
[nal.core :refer :all]))
|
||||
|
||||
(deftest test-inference
|
||||
(are [a1 a2] (= (set a1) (set (apply inference a2)))
|
||||
'({:statement [==>
|
||||
[&| [--> [ext-set tim] [int-set driving]]]
|
||||
[--> [ext-set tim] [int-set dead]]]
|
||||
:truth [1.0 0.81]
|
||||
:task-type :judgement
|
||||
|
||||
'({:statement [==>
|
||||
[&| [--> [ext-set tim] [int-set driving]]]
|
||||
[--> [ext-set tim] [int-set dead]]]
|
||||
:truth [1.0 0.81]
|
||||
:task-type :judgement
|
||||
:occurrence 1})
|
||||
|
||||
'[{:statement [--> [ext-set tim] [int-set drunk]]
|
||||
:truth [1 0.9]
|
||||
:task-type :judgement
|
||||
'[{:statement [--> [ext-set tim] [int-set drunk]]
|
||||
:truth [1 0.9]
|
||||
:task-type :judgement
|
||||
:occurrence 1}
|
||||
|
||||
{:statement [==> [&| [--> [ind-var X] [int-set drunk]] [--> [ind-var X] [int-set driving]]]
|
||||
[--> [ind-var X] [int-set dead]]]
|
||||
:truth [1 0.9]
|
||||
{:statement [==> [&| [--> [ind-var X] [int-set drunk]] [--> [ind-var X] [int-set driving]]]
|
||||
[--> [ind-var X] [int-set dead]]]
|
||||
:truth [1 0.9]
|
||||
:occurrence 0}]
|
||||
|
||||
'({:occurrence 1
|
||||
@@ -50,38 +47,38 @@
|
||||
:task-type :judgement
|
||||
:truth [1
|
||||
0.44751381215469616]})
|
||||
'[{:statement [--> [* a1 a2 a3] m]
|
||||
:truth [1 0.9]
|
||||
:task-type :judgement
|
||||
'[{:statement [--> [* a1 a2 a3] m]
|
||||
:truth [1 0.9]
|
||||
:task-type :judgement
|
||||
:occurrence 1}
|
||||
|
||||
{:statement a1
|
||||
:truth [1 0.9]
|
||||
{:statement a1
|
||||
:truth [1 0.9]
|
||||
:occurrence 0}]
|
||||
|
||||
'[{:statement [=|> [--> [* a1 a2 a3] m] a1],
|
||||
:task-type :judgement,
|
||||
'[{:statement [=|> [--> [* a1 a2 a3] m] a1],
|
||||
:task-type :judgement,
|
||||
:occurrence 1,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [<|> a1 [--> [* a1 a2 a3] m]],
|
||||
:task-type :judgement,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [<|> a1 [--> [* a1 a2 a3] m]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1,
|
||||
:truth [1.0 0.44751381215469616]}
|
||||
{:statement [&| [--> [* a1 a2 a3] m] a1],
|
||||
:task-type :judgement,
|
||||
:truth [1.0 0.44751381215469616]}
|
||||
{:statement [&| [--> [* a1 a2 a3] m] a1],
|
||||
:task-type :judgement,
|
||||
:occurrence 1,
|
||||
:truth [1.0 0.81]}
|
||||
{:statement [=|> a1 [--> [* a1 a2 a3] m]],
|
||||
:task-type :judgement,
|
||||
:truth [1.0 0.81]}
|
||||
{:statement [=|> a1 [--> [* a1 a2 a3] m]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1,
|
||||
:truth [1 0.44751381215469616]}]
|
||||
'[{:statement a1
|
||||
:truth [1 0.9]
|
||||
:task-type :judgement
|
||||
:truth [1 0.44751381215469616]}]
|
||||
'[{:statement a1
|
||||
:truth [1 0.9]
|
||||
:task-type :judgement
|
||||
:occurrence 1}
|
||||
|
||||
{:statement [--> [* a1 a2 a3] m]
|
||||
:truth [1 0.9]
|
||||
{:statement [--> [* a1 a2 a3] m]
|
||||
:truth [1 0.9]
|
||||
:occurrence 0}]
|
||||
|
||||
'({:occurrence 1
|
||||
@@ -106,166 +103,56 @@
|
||||
:statement a1
|
||||
:task-type :judgement
|
||||
:truth [1 0.44751381215469616]})
|
||||
'[{:statement [conj a1 a2 a3]
|
||||
:truth [1 0.9]
|
||||
:task-type :judgement
|
||||
'[{:statement [conj a1 a2 a3]
|
||||
:truth [1 0.9]
|
||||
:task-type :judgement
|
||||
:occurrence 1}
|
||||
|
||||
{:statement a1
|
||||
:truth [1 0.9]
|
||||
{:statement a1
|
||||
:truth [1 0.9]
|
||||
:occurrence 0}]
|
||||
|
||||
'[{:statement [=|> [--> M S] [[--> M S] [--> M P]]],
|
||||
:task-type :judgement,
|
||||
'[{:statement [=|> [--> M S] [[--> M S] [--> M P]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [<|> [--> M S] [[--> M S] [--> M P]]],
|
||||
:task-type :judgement,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [<|> [--> M S] [[--> M S] [--> M P]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1,
|
||||
:truth [1.0 0.44751381215469616]}
|
||||
{:statement [&| [--> M S] [[--> M S] [--> M P]]],
|
||||
:task-type :judgement,
|
||||
:truth [1.0 0.44751381215469616]}
|
||||
{:statement [&| [--> M S] [[--> M S] [--> M P]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1,
|
||||
:truth [1.0 0.81]}
|
||||
{:statement [=|> [[--> M S] [--> M P]] [--> M S]],
|
||||
:task-type :judgement,
|
||||
:truth [1.0 0.81]}
|
||||
{:statement [=|> [[--> M S] [--> M P]] [--> M S]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1,
|
||||
:truth [1 0.44751381215469616]}]
|
||||
'[{:statement [[--> M S] [--> M P]]
|
||||
:truth [1 0.9]
|
||||
:task-type :judgement
|
||||
:truth [1 0.44751381215469616]}]
|
||||
'[{:statement [[--> M S] [--> M P]]
|
||||
:truth [1 0.9]
|
||||
:task-type :judgement
|
||||
:occurrence 1}
|
||||
|
||||
{:statement [--> M S]
|
||||
:truth [1 0.9]
|
||||
{:statement [--> M S]
|
||||
:truth [1 0.9]
|
||||
:occurrence 0}]
|
||||
|
||||
'({:statement [conj [--> [dep-var Y] [int-set B]]
|
||||
[==> [--> [ext-set A] [int-set Y]] [--> [dep-var Y] P]]]
|
||||
:truth [1.0 0.81]
|
||||
:task-type :judgement
|
||||
'({:statement [conj [--> [dep-var Y] [int-set B]]
|
||||
[==> [--> [ext-set A] [int-set Y]] [--> [dep-var Y] P]]]
|
||||
:truth [1.0 0.81]
|
||||
:task-type :judgement
|
||||
:occurrence 1}
|
||||
{:statement [==>
|
||||
[conj [--> [ext-set A] [int-set Y]] [--> [ind-var X] [int-set B]]]
|
||||
[--> [ind-var X] P]]
|
||||
:truth [1 0.44751381215469616]
|
||||
:task-type :judgement
|
||||
{:statement [==>
|
||||
[conj [--> [ext-set A] [int-set Y]] [--> [ind-var X] [int-set B]]]
|
||||
[--> [ind-var X] P]]
|
||||
:truth [1 0.44751381215469616]
|
||||
:task-type :judgement
|
||||
:occurrence 1})
|
||||
'[{:statement [==> [--> [ext-set A] [int-set Y]] [--> [ext-set A] P]]
|
||||
:truth [1 0.9]
|
||||
:task-type :judgement
|
||||
'[{:statement [==> [--> [ext-set A] [int-set Y]] [--> [ext-set A] P]]
|
||||
:truth [1 0.9]
|
||||
:task-type :judgement
|
||||
:occurrence 1}
|
||||
|
||||
{:statement [--> [ext-set A] [int-set B]]
|
||||
:truth [1 0.9]
|
||||
:occurrence 1}]
|
||||
'({:statement [</>
|
||||
[seq-conj [--> chess competition] [:interval 1000]]
|
||||
[--> sport competition]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.44751381215469616]}
|
||||
{:statement [seq-conj
|
||||
[--> chess competition]
|
||||
[:interval 1000]
|
||||
[--> sport competition]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.81]}
|
||||
{:statement [pred-impl
|
||||
[seq-conj [--> chess competition] [:interval 1000]]
|
||||
[--> sport competition]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [retro-impl
|
||||
[--> sport competition]
|
||||
[seq-conj [--> chess competition] [:interval 1000]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [--> sport chess],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [--> chess sport],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [<=> [--> chess [ind-var X]] [--> sport [ind-var X]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.44751381215469616]}
|
||||
{:statement [conj [--> chess [dep-var Y]] [--> sport [dep-var Y]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.81]}
|
||||
{:statement [<-> sport chess],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.44751381215469616]}
|
||||
{:statement [==> [--> chess [ind-var X]] [--> sport [ind-var X]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [==> [--> chess [ind-var X]] [--> sport [ind-var X]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [==> [--> sport [ind-var X]] [--> chess [ind-var X]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [==> [--> sport [ind-var X]] [--> chess [ind-var X]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [--> [int-dif chess sport] competition],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [0.0 0.81]}
|
||||
{:statement [--> [| chess sport] competition],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.81]}
|
||||
{:statement [--> [int-dif sport chess] competition],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [0.0 0.81]}
|
||||
{:statement [--> [ext-inter chess sport] competition],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.81]}
|
||||
{:statement [pred-impl
|
||||
[seq-conj [--> chess [ind-var X]] [:interval 1000]]
|
||||
[--> sport [ind-var X]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [</>
|
||||
[seq-conj [--> chess [ind-var X]] [:interval 1000]]
|
||||
[--> sport [ind-var X]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.44751381215469616]}
|
||||
{:statement [retro-impl
|
||||
[--> sport [ind-var X]]
|
||||
[seq-conj [--> chess [ind-var X]] [:interval 1000]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [seq-conj
|
||||
[--> chess [dep-var Y]]
|
||||
[:interval 1000]
|
||||
[--> sport [dep-var Y]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.81]})
|
||||
['{:statement [--> sport competition]
|
||||
:truth [1 0.9]
|
||||
:task-type :judgement
|
||||
:occurrence 1000}
|
||||
|
||||
'{:statement [--> chess competition]
|
||||
:truth [1 0.9]
|
||||
:occurrence 0}]))
|
||||
{:statement [--> [ext-set A] [int-set B]]
|
||||
:truth [1 0.9]
|
||||
:occurrence 1}]))
|
||||
|
||||
+105
-10
@@ -1,21 +1,116 @@
|
||||
(ns nal.test.deriver
|
||||
(:require [clojure.test :refer :all]
|
||||
[nal.deriver :refer :all]
|
||||
[nal.deriver.rules :refer [compile-rules]]
|
||||
[nal.rules :as r]))
|
||||
|
||||
(def rules
|
||||
(compile-rules '([(P --> M) (S --> M) |- (S <-> P)
|
||||
:post (:t/comparison :d/weak :allow-backward)
|
||||
:pre ((:!= S P))])))
|
||||
(def result
|
||||
'({:statement [</>
|
||||
[seq-conj [--> chess competition] [:interval 1000]]
|
||||
[--> sport competition]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.44751381215469616]}
|
||||
{:statement [seq-conj
|
||||
[--> chess competition]
|
||||
[:interval 1000]
|
||||
[--> sport competition]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.81]}
|
||||
{:statement [pred-impl
|
||||
[seq-conj [--> chess competition] [:interval 1000]]
|
||||
[--> sport competition]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [retro-impl
|
||||
[--> sport competition]
|
||||
[seq-conj [--> chess competition] [:interval 1000]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [--> sport chess],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [--> chess sport],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [<=> [--> chess [ind-var X]] [--> sport [ind-var X]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.44751381215469616]}
|
||||
{:statement [conj [--> chess [dep-var Y]] [--> sport [dep-var Y]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.81]}
|
||||
{:statement [<-> sport chess],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.44751381215469616]}
|
||||
{:statement [==> [--> chess [ind-var X]] [--> sport [ind-var X]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [==> [--> chess [ind-var X]] [--> sport [ind-var X]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [==> [--> sport [ind-var X]] [--> chess [ind-var X]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [==> [--> sport [ind-var X]] [--> chess [ind-var X]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [--> [int-dif chess sport] competition],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [0.0 0.81]}
|
||||
{:statement [--> [| chess sport] competition],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.81]}
|
||||
{:statement [--> [int-dif sport chess] competition],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [0.0 0.81]}
|
||||
{:statement [--> [ext-inter chess sport] competition],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.81]}
|
||||
{:statement [pred-impl
|
||||
[seq-conj [--> chess [ind-var X]] [:interval 1000]]
|
||||
[--> sport [ind-var X]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [</>
|
||||
[seq-conj [--> chess [ind-var X]] [:interval 1000]]
|
||||
[--> sport [ind-var X]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.44751381215469616]}
|
||||
{:statement [retro-impl
|
||||
[--> sport [ind-var X]]
|
||||
[seq-conj [--> chess [ind-var X]] [:interval 1000]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1 0.44751381215469616]}
|
||||
{:statement [seq-conj
|
||||
[--> chess [dep-var Y]]
|
||||
[:interval 1000]
|
||||
[--> sport [dep-var Y]]],
|
||||
:task-type :judgement,
|
||||
:occurrence 1000,
|
||||
:truth [1.0 0.81]}))
|
||||
|
||||
(deftest test-generate-conclusions
|
||||
(is (= (set [{:occurrence 1000
|
||||
:statement '[<-> sport chess]
|
||||
:task-type :judgement
|
||||
:truth [1.0 0.44751381215469616]}])
|
||||
(is (= (set result)
|
||||
(set (generate-conclusions
|
||||
(rules :judgement)
|
||||
(r/rules :judgement)
|
||||
'{:statement [--> sport competition]
|
||||
:truth [1 0.9]
|
||||
:task-type :judgement
|
||||
|
||||
@@ -1,88 +0,0 @@
|
||||
(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)))))
|
||||
@@ -1,156 +0,0 @@
|
||||
(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)))))
|
||||
@@ -1,161 +0,0 @@
|
||||
(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)))))
|
||||
@@ -1,114 +0,0 @@
|
||||
(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)))))
|
||||
@@ -1,123 +0,0 @@
|
||||
(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)))))
|
||||
@@ -1,91 +0,0 @@
|
||||
(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)))))
|
||||
@@ -1,42 +0,0 @@
|
||||
(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)))
|
||||
@@ -1,67 +0,0 @@
|
||||
(ns narjure.test.memory.redis
|
||||
(:require [clojure.test :refer :all]
|
||||
[narjure.memory.redis :as r]
|
||||
[narjure.memory.api :as m]
|
||||
[taoensso.carmine :as c]))
|
||||
|
||||
(def config
|
||||
{:pool {}
|
||||
:spec {:host "127.0.0.1" :port 6379}})
|
||||
|
||||
(def mem (r/->RedisMemory config))
|
||||
|
||||
(def concept1 '[--> tim cat])
|
||||
(def concept2 'tim)
|
||||
|
||||
(def belief1 {:frequency (float 0.9)
|
||||
:confidence (float 0.2)
|
||||
:occurrence 1
|
||||
:term concept1})
|
||||
|
||||
(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)
|
||||
:durability (float 0.1)
|
||||
: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))))
|
||||
|
||||
(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 (= [goal1]
|
||||
(map :task (m/goals mem concept1)))))
|
||||
|
||||
(is (:task (first (m/beliefs mem concept1))))
|
||||
|
||||
(m/add-term mem concept2)
|
||||
(m/add-termlink mem concept1 (termlink concept2))
|
||||
(is (= [(termlink concept2)]
|
||||
(map #(dissoc % :id) (m/termlinks mem concept1))))
|
||||
(is (= concept2 (m/term mem (:concept (first (m/termlinks mem concept1)))))))
|
||||
Reference in New Issue
Block a user