Skip to content

Job only running one time #4

Description

@naglalakk

I'm wondering why my queue is not getting processed. When I run this code a single job is processed but then if I add debugging I see that the queue is just filling up. My code

main = do
	hworker <- create "notifier" (State state)
	forkIO (worker hworker)
	forkIO (monitor hworker)
	forever $ queue hworker (ServiceAction "test" "testing" 1) >> threadDelay 5000000

This should push a new job into the queue every 5 seconds, which it does. But it's not processing it. I checked the failed jobs queue and it's completely empty so they are not failing.

This is the Job instance

instance Job State ServiceJob where
	job (State mvar) (ServiceAction service action id) = do 
		clients <- takeMVar mvar
		let message = "service:" 
				   ++ (service) 
				   ++ "|action:" 
				   ++ (action) 
				   ++ "|id:" 
				   ++ show id
		print message >> return Success
		--broadcast (T.pack message) clients
		-- return Success
	job (State mvar) (SearchIndex service id) = do
		putStrLn "search"
		return Success

Activity

  1. naglalakk commented on Dec 13, 2017

    @naglalakk
    Author

    UPDATES

    I'm running this now with debugging. I started two workers, each running a seperate job. Now they both run once and then the second worker goes into a really weird non stop loop. Note: the loop is seen in the second snippet here below (redis debugger)

    Hworker debugger:

    
    *Main> ("notifier",("WORKER RUNNING","[\"450733c6-3c6f-451c-9a25-9a5570b19f85\",{\"tag\":\"ServiceAction\",\"contents\":[\"test\",\"testing\",1]}]"))
    "service:test|action:testing|id:1"
    ("notifier",("JOB COMPLETE","[\"450733c6-3c6f-451c-9a25-9a5570b19f85\",{\"tag\":\"ServiceAction\",\"contents\":[\"test\",\"testing\",1]}]"))
    ("notifier",("MONITOR RV",1))
    ("notifier",("MONITOR REQUEUED","[\"df925ed6-4c86-4989-b967-99230d13715a\",{\"tag\":\"ServiceAction\",\"contents\":[\"test\",\"testing\",1]}]"))
    ("notifier",("WORKER RUNNING","[\"df925ed6-4c86-4989-b967-99230d13715a\",{\"tag\":\"ServiceAction\",\"contents\":[\"test\",\"testing\",1]}]"))
    ("search",("WORKER RUNNING","[\"38410142-6f83-4348-8f95-d2ab1dd19570\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]"))
    ("search",("JOB COMPLETE","[\"38410142-6f83-4348-8f95-d2ab1dd19570\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]"))
    ("notifier",("MONITOR RV",1))
    ("notifier",("MONITOR REQUEUED","[\"df925ed6-4c86-4989-b967-99230d13715a\",{\"tag\":\"ServiceAction\",\"contents\":[\"test\",\"testing\",1]}]"))
    

    Redis debug

    1513173783.022352 [0 172.17.0.1:46054] "EVAL" "local job = redis.call('rpop',KEYS[1])\nif job ~= nil then\n  redis.call('hset', KEYS[2], job, ARGV[1])\n  return job\nelse\n  return nil\nend" "2" "hworker-jobs-search" "hworker-progress-search" "\"2017-12-13T14:03:03.021753Z\""
    1513173783.022405 [0 lua] "rpop" "hworker-jobs-search"
    1513173783.034465 [0 172.17.0.1:46054] "EVAL" "local job = redis.call('rpop',KEYS[1])\nif job ~= nil then\n  redis.call('hset', KEYS[2], job, ARGV[1])\n  return job\nelse\n  return nil\nend" "2" "hworker-jobs-search" "hworker-progress-search" "\"2017-12-13T14:03:03.03398Z\""
    1513173783.034511 [0 lua] "rpop" "hworker-jobs-search"
    1513173783.046367 [0 172.17.0.1:46054] "EVAL" "local job = redis.call('rpop',KEYS[1])\nif job ~= nil then\n  redis.call('hset', KEYS[2], job, ARGV[1])\n  return job\nelse\n  return nil\nend" "2" "hworker-jobs-search" "hworker-progress-search" "\"2017-12-13T14:03:03.045846Z\""
    1513173783.046417 [0 lua] "rpop" "hworker-jobs-search"
    1513173783.057326 [0 172.17.0.1:46054] "EVAL" "local job = redis.call('rpop',KEYS[1])\nif job ~= nil then\n  redis.call('hset', KEYS[2], job, ARGV[1])\n  return job\nelse\n  return nil\nend" "2" "hworker-jobs-search" "hworker-progress-search" "\"2017-12-13T14:03:03.056814Z\""
    1513173783.057373 [0 lua] "rpop" "hworker-jobs-search"
    1513173783.070731 [0 172.17.0.1:46054] "EVAL" "local job = redis.call('rpop',KEYS[1])\nif job ~= nil then\n  redis.call('hset', KEYS[2], job, ARGV[1])\n  return job\nelse\n  return nil\nend" "2" "hworker-jobs-search" "hworker-progress-search" "\"2017-12-13T14:03:03.069905Z\""
    1513173783.070801 [0 lua] "rpop" "hworker-jobs-search"
    
  2. dbp commented on Dec 13, 2017

    @dbp
    Contributor
  3. naglalakk commented on Dec 13, 2017

    @naglalakk
    Author

    Hi @dbp thank you for answering. I'm trying to see what might be wrong here. I did name my queues differently. Now I have this

    	let config_1 = (defaultHworkerConfig "notifier" (State state)) { hwconfigDebug = True}
    	let config_2 = (defaultHworkerConfig "search" (State state)) { hwconfigDebug = True}
    
    	-- Start the worker
    	nworker <- createWith config_1
    	sworker <- createWith config_2
    
    	forkIO (monitor nworker)
    	forkIO (worker nworker)
    	forkIO (worker sworker)
    	forkIO (monitor sworker)
    	forkIO (forever $ queue nworker (ServiceAction "test" "testing" 1) >> threadDelay 5000000)
    	forkIO (forever $ queue sworker (SearchIndex "test" 2) >> threadDelay 5000000)
    

    On my first run after changing this, the queue named search worked but on a second run it didn't work anymore. I have the debugger running and occasionally the monitor says "MONITOR REQUEUED" but then nothing. I'm writing the results of both jobs to a file so I should see a new line of text when it runs.

    Could this be because of the threads I create? I'm very new to threading in Haskell

    The redis debugger shows that it's definitely pushing the tasks into the queue but does not seem to process them.

    
    1513180391.971306 [0 172.17.0.1:46532] "HKEYS" "hworker-progress-search"
    1513180393.126282 [0 172.17.0.1:46532] "LPUSH" "hworker-jobs-search" "[\"fc3f116c-4ba3-46b6-b233-7f6fca9512d3\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]"
    1513180398.130543 [0 172.17.0.1:46532] "LPUSH" "hworker-jobs-search" "[\"8061a508-6a84-4fe1-b810-dbbd51d2d3f4\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]"
    1513180403.136544 [0 172.17.0.1:46532] "LPUSH" "hworker-jobs-search" "[\"997a9d6d-fa3a-427b-9d18-b94eda986cfc\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]"
    1513180403.974982 [0 172.17.0.1:46532] "HKEYS" "hworker-progress-search"
    1513180408.137567 [0 172.17.0.1:46532] "LPUSH" "hworker-jobs-search" "[\"d8af7c59-9357-4d0f-b80a-4694e6e248c3\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]"
    1513180413.141071 [0 172.17.0.1:46532] "LPUSH" "hworker-jobs-search" "[\"de15be50-058f-432e-97cd-d7434f2197f8\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]"
    1513180415.978862 [0 172.17.0.1:46532] "HKEYS" "hworker-progress-search"
    1513180418.145018 [0 172.17.0.1:46532] "LPUSH" "hworker-jobs-search" "[\"f9dd7014-8c50-419b-83d1-9ec5f9359e17\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]"
    1513180423.150860 [0 172.17.0.1:46532] "LPUSH" "hworker-jobs-search" "[\"be5233c4-b648-48cd-85bf-7b18ca1699d9\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]"
    1513180427.982957 [0 172.17.0.1:46532] "HKEYS" "hworker-progress-search"
    1513180428.155229 [0 172.17.0.1:46532] "LPUSH" "hworker-jobs-search" "[\"181b57ce-eedd-4e94-9ea7-f0c1f8ab002b\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]"
    
    
  4. naglalakk commented on Dec 13, 2017

    @naglalakk
    Author

    This is what makes me think it has something to do with threading. The debugger is showing garbled output

    *Main> ("search",("W(O"RnKoEtRi fRiUeNrN"I,N(G""W,O"R[K\E"Re dR8UfN0N9IdN4G-"3,2"a[3\-"41ba5c6f-3a761331--6382f0a0-f4439188b-b8727b\8"-,9{b\5"6t8a3ga\c"3:f\b"bS\e"r,v{i\c"etAacgt\i"o:n\\""S,e\r"vciocnetAecnttiso\n"\:"[,\\""tceosntt\e"n,t\s"\t"e:s[t\i"ntge\s"t,\1"],}\]""t)e)s
    ting\",1]}(]""s)e)a
    rch",("JOB COMPLETE","[\"ed8f09d4-32a3-4b56-a631-3200f498bb77\",{\"tag\":\"ServiceAction\",\"contents\":[\"test\",\"testing\",1]}]"))
    ("search",("WORKER RUNNING","[\"cab85ce4-ba73-400a-8f86-22144ac0fca7\",{\"tag\":\"ServiceAction\",\"contents\":[\"test\",\"testing\",1]}]"))
    
  5. dbp commented on Dec 13, 2017

    @dbp
    Contributor

    So based on the program fragment you gave above, one problem might be that every thread you forked off with forkIO, which means that the main thread will eventually terminate, which will then kill the program. e.g., the following program will print (on my machine) about 1200 times and then quit:

    import Control.Concurrent (forkIO)
    import Control.Monad (forever)
    
    main :: IO ()
    main = do forkIO $ forever (putStrLn "Echo")
              return ()
    

    Usually the way around this is to leave one of them (usually the most important one, in this case perhaps the monitor?) in the main thread (i.e., just leave off the forkIO call). If that's not it, is there a way you can post a little more code?

  6. naglalakk commented on Dec 13, 2017

    @naglalakk
    Author

    I tried this but I get the same results. Here is my entire code. Don't mind the Websocket stuff it's not really being used at the moment

    {-# LANGUAGE MultiParamTypeClasses #-}
    {-# LANGUAGE DeriveGeneric         #-}
    {-# LANGUAGE OverloadedStrings     #-}
    
    module Main where
    
    import Data.List			           (intercalate)
    import GHC.Generics			   (Generic)
    import Data.Aeson			           (FromJSON, ToJSON)
    import Control.Concurrent.MVar     (MVar, newMVar, putMVar, takeMVar)
    import Control.Concurrent               (forkIO, threadDelay)
    import Control.Monad                      (forever, forM_)
    import System.Hworker
    import qualified Data.Text as T
    import qualified Network.WebSockets as WS
    
    
    type Client = (T.Text, WS.Connection)
    type ServerState = [Client]
    
    data ServiceJob = 
    	  ServiceAction String String Int 
    	| SearchIndex   String Int deriving (Generic, Show)
    
    data State = State (MVar ServerState)
    
    instance ToJSON ServiceJob
    instance FromJSON ServiceJob
    
    instance Job State ServiceJob where
    	job (State mvar) (ServiceAction service action id) = do 
    		clients <- takeMVar mvar
    		let message = "service:" 
    				   ++ (service) 
    				   ++ "|action:" 
    				   ++ (action) 
    				   ++ "|id:" 
    				   ++ show id
    		appendFile "/Users/Donna/Code/K0TT/repos/cherry-service-worker-test/test.txt" (message ++ "\n")
    		>> return Success
    	job (State mvar) (SearchIndex service id) = do
    		appendFile "/Users/Donna/Code/K0TT/repos/cherry-service-worker-test/test.txt" "search\n"
    		>> return Success
    
    broadcast :: T.Text -> ServerState -> IO ()
    broadcast message clients = do
    	forM_ clients $ \(_, conn) -> WS.sendTextData conn message
    
    newServerState :: ServerState
    newServerState = []
    
    main = do
    	state <- newMVar newServerState
    
    	let config_1 = (defaultHworkerConfig "notifier" (State state)) { hwconfigDebug = True}
    	let config_2 = (defaultHworkerConfig "search" (State state)) { hwconfigDebug = True}
    
    	nworker <- createWith config_1
    	sworker <- createWith config_2
    
    	-- Start the notifier worker
    	forkIO (worker nworker)
    
    	-- Start the search worker 
    	forkIO (worker sworker)
    	forkIO (monitor sworker)
    
    	-- Add new ServiceAction every 5 secs
    	forkIO (forever $ queue nworker (ServiceAction "test" "testing" 1) >> threadDelay 5000000)
    
    	-- Add new SearchIndex every 10 secs
    	forkIO (forever $ queue sworker (SearchIndex "test" 2) >> threadDelay 10000000)
    
    	-- Monitor the notifier worker
    	monitor nworker
    
    
  7. dbp commented on Dec 13, 2017

    @dbp
    Contributor

    So there are two things I notice:

    1. Your ServiceAction takes the MVar but never refills it, so this will lock the thread the second time it runs (MVars are either empty or full, and different functions have different blocking behavior -- takeMVar waits until it gets filled if it is empty, but since that'll never happen here, it'll wait forever).

    2. It shouldn't matter here, but you don't actually need two separate queues (if you did that intentionally, ignore me!). A given queue can store a single data type, but it can have multiple variants (in this case, ServiceActions and SearchIndexes). What you are doing instead is creating two different queues, both of which can store either type of job, and then only storing one type in each queue. If you wanted to separate them completely (because in the future you might want the ServiceActions to run on one server and the SearchIndexes on another), it might be better to create different data types for them, so that way you couldn't accidentally queue the wrong type of job into a queue.

    I ran your program locally (with minor changes for the path of where it writes the log messages), and I see the behavior I expect based on the above -- the ServiceAction task runs exactly once (as it then gets blocked), and the SearchIndex runs every ten seconds. You are using redis 2.8 right? We really need to support newer versions, but they changed some stuff in EVAL which we rely heavily on!

  8. naglalakk commented on Dec 13, 2017

    @naglalakk
    Author

    Ahh of course, that was it. I was not paying attention to the MVar because I wasn't using it yet. Thank you so much and sorry to waste your time with this. And you are correct, I don't need two queues but I thought that maybe that was the reason for the delay of the queues. Brilliant, now everything works. Thanks again 👍

  9. dbp commented on Dec 13, 2017

    @dbp
    Contributor

    Glad it's working, and no worries -- happy I was able to help!

  10. added a commit that references this issue on Dec 30, 2024
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions