mirror of
https://github.com/simplex-chat/simplexmq.git
synced 2026-07-28 20:50:10 +00:00
smp server: fixed logging format for journal store errors
This commit is contained in:
@@ -306,7 +306,7 @@ instance MsgStoreClass JournalMsgStore where
|
||||
(msg :) <$> getMsg msgs hs
|
||||
|
||||
writeMsg :: JournalMsgStore -> JournalMsgQueue -> Bool -> Message -> ExceptT ErrorType IO (Maybe (Message, Bool))
|
||||
writeMsg ms q@JournalMsgQueue {queue = JMQueue {queueDirectory, queueLock, statePath}, handles} logState !msg =
|
||||
writeMsg ms q@JournalMsgQueue {queue = JMQueue {queueDirectory, queueLock, statePath}, handles} logState msg =
|
||||
tryStore "writeMsg" $ withLock' queueLock "writeMsg" $ do
|
||||
st@MsgQueueState {canWrite, size} <- readTVarIO (state q)
|
||||
let empty = size == 0
|
||||
@@ -319,8 +319,8 @@ instance MsgStoreClass JournalMsgStore where
|
||||
else pure Nothing
|
||||
where
|
||||
JournalStoreConfig {quota, maxMsgCount} = config ms
|
||||
!msgQuota = MessageQuota {msgId = msgId msg, msgTs = msgTs msg}
|
||||
writeToJournal st@MsgQueueState {writeState, readState = rs, size} canWrt' msg' = do
|
||||
msgQuota = MessageQuota {msgId = msgId msg, msgTs = msgTs msg}
|
||||
writeToJournal st@MsgQueueState {writeState, readState = rs, size} canWrt' !msg' = do
|
||||
let msgStr = strEncode msg' `B.snoc` '\n'
|
||||
msgLen = fromIntegral $ B.length msgStr
|
||||
hs <- maybe createQueueDir pure =<< readTVarIO handles
|
||||
@@ -384,7 +384,7 @@ tryStore op a =
|
||||
ExceptT $
|
||||
(Right <$> a) `catchAny` \e ->
|
||||
let e' = op <> " " <> show e
|
||||
in logError ("STORE ERROR " <> T.pack e') $> Left (STORE e')
|
||||
in logError ("STORE " <> T.pack e') $> Left (STORE e')
|
||||
|
||||
openMsgQueue :: JournalMsgStore -> JMQueue -> IO JournalMsgQueue
|
||||
openMsgQueue ms q@JMQueue {queueDirectory = dir, statePath} = do
|
||||
@@ -420,7 +420,7 @@ chooseReadJournal q log' hs = do
|
||||
|
||||
updateQueueState :: JournalMsgQueue -> Bool -> MsgQueueHandles -> MsgQueueState -> IO ()
|
||||
updateQueueState q log' hs st = do
|
||||
unless (validQueueState st) $ E.throwIO $ userError $ "updating to invalid state: " <> show st
|
||||
unless (validQueueState st) $ E.throwIO $ userError $ "updateQueueState, updating to invalid state, " <> show st
|
||||
when log' $ appendState (stateHandle hs) st
|
||||
atomically $ writeTVar (state q) st
|
||||
|
||||
@@ -465,7 +465,7 @@ openJournals dir st@MsgQueueState {readState = rs, writeState = ws} = do
|
||||
wjId = journalId ws
|
||||
openJournal rs >>= \case
|
||||
Left path -> do
|
||||
logError $ "STORE ERROR no read file " <> T.pack path <> ", creating new file"
|
||||
logError $ "STORE openJournals, no read file " <> T.pack path <> ", creating new file"
|
||||
rh <- createNewJournal dir rjId
|
||||
let st' = newMsgQueueState rjId
|
||||
pure (st', rh, Nothing)
|
||||
@@ -477,7 +477,7 @@ openJournals dir st@MsgQueueState {readState = rs, writeState = ws} = do
|
||||
fixFileSize rh $ byteCount rs
|
||||
openJournal ws >>= \case
|
||||
Left path -> do
|
||||
logError $ "STORE ERROR no write file " <> T.pack path <> ", creating new file"
|
||||
logError $ "STORE openJournals, no write file " <> T.pack path <> ", creating new file"
|
||||
wh <- createNewJournal dir wjId
|
||||
let size' = msgCount rs - msgPos rs
|
||||
st' = st {writeState = newJournalState wjId, size = size'} -- we don't amend canWrite to trigger QCONT
|
||||
@@ -498,17 +498,17 @@ fixFileSize h pos = do
|
||||
size <- IO.hFileSize h
|
||||
if
|
||||
| size > pos' -> do
|
||||
logWarn $ "STORE WARNING truncating file size from " <> tshow size <> " to " <> tshow pos
|
||||
logWarn $ "STORE fixFileSize, truncating from " <> tshow size <> " to " <> tshow pos
|
||||
IO.hSetFileSize h pos'
|
||||
| size < pos' ->
|
||||
-- From code logic this can't happen.
|
||||
E.throwIO $ userError $ "file size " <> show size <> " is smaller than position " <> show pos
|
||||
E.throwIO $ userError $ "fixFileSize, file size " <> show size <> " is smaller than position " <> show pos
|
||||
| otherwise -> pure ()
|
||||
|
||||
removeJournal :: FilePath -> JournalState t -> IO ()
|
||||
removeJournal dir JournalState {journalId} = do
|
||||
let path = journalFilePath dir journalId
|
||||
removeFile path `catchAny` (\e -> logError $ "STORE ERROR removing file " <> T.pack path <> ": " <> tshow e)
|
||||
removeFile path `catchAny` (\e -> logError $ "STORE removeJournal, " <> T.pack path <> ", " <> tshow e)
|
||||
|
||||
-- This function is supposed to be resilient to crashes while updating state files,
|
||||
-- and also resilient to crashes during its execution.
|
||||
@@ -526,10 +526,10 @@ readWriteQueueState JournalMsgStore {random, config} statePath =
|
||||
[] -> writeNewQueueState
|
||||
_ -> do
|
||||
r@(st, _) <- useLastLine (length ls) True ls
|
||||
unless (validQueueState st) $ E.throwIO $ userError $ "read invalid invalid: " <> show st
|
||||
unless (validQueueState st) $ E.throwIO $ userError $ "readWriteQueueState, read invalid invalid, " <> show st
|
||||
pure r
|
||||
writeNewQueueState = do
|
||||
logWarn $ "STORE WARNING: empty queue state in " <> T.pack statePath <> ", initialized"
|
||||
logWarn $ "STORE readWriteQueueState, empty queue state in " <> T.pack statePath <> ", initialized"
|
||||
st <- newMsgQueueState <$> newJournalId random
|
||||
writeQueueState st
|
||||
useLastLine len isLastLine ls = case strDecode $ LB.toStrict $ last ls of
|
||||
@@ -543,13 +543,13 @@ readWriteQueueState JournalMsgStore {random, config} statePath =
|
||||
Left e -- if the last line failed to parse
|
||||
| isLastLine -> case init ls of -- or use the previous line
|
||||
[] -> do
|
||||
logWarn $ "STORE WARNING: invalid 1-line queue state " <> T.pack statePath <> ", initialized"
|
||||
logWarn $ "STORE readWriteQueueState, invalid 1-line queue state " <> T.pack statePath <> ", initialized"
|
||||
st <- newMsgQueueState <$> newJournalId random
|
||||
backupWriteQueueState st
|
||||
ls' -> do
|
||||
logWarn $ "STORE WARNING: invalid last line in queue state " <> T.pack statePath <> ", using the previous line"
|
||||
logWarn $ "STORE readWriteQueueState, invalid last line in queue state " <> T.pack statePath <> ", using the previous line"
|
||||
useLastLine len False ls'
|
||||
| otherwise -> E.throwIO $ userError $ "reading queue state " <> statePath <> ": " <> show e
|
||||
| otherwise -> E.throwIO $ userError $ "readWriteQueueState, reading queue state " <> statePath <> ", " <> show e
|
||||
backupWriteQueueState st = do
|
||||
-- State backup is made in two steps to mitigate the crash during the backup.
|
||||
-- Temporary backup file will be used when it is present.
|
||||
@@ -598,7 +598,7 @@ closeMsgQueue q = readTVarIO (handles q) >>= mapM_ closeHandles
|
||||
removeQueueDirectory :: JournalMsgStore -> RecipientId -> IO ()
|
||||
removeQueueDirectory st rId =
|
||||
let dir = msgQueueDirectory st rId
|
||||
in removePathForcibly dir `catchAny` (\e -> logError $ "STORE ERROR removeQueueDirectory " <> T.pack dir <> ": " <> tshow e)
|
||||
in removePathForcibly dir `catchAny` (\e -> logError $ "STORE removeQueueDirectory, " <> T.pack dir <> ", " <> tshow e)
|
||||
|
||||
hAppend :: Handle -> Int64 -> ByteString -> IO ()
|
||||
hAppend h pos s = do
|
||||
@@ -614,7 +614,7 @@ hGetMsgAt h pos = do
|
||||
Right !msg ->
|
||||
let !len = fromIntegral (B.length s) + 1
|
||||
in pure (msg, len)
|
||||
Left e -> E.throwIO $ userError $ "Error parsing message: " <> e
|
||||
Left e -> E.throwIO $ userError $ "hGetMsgAt, error parsing message, " <> e
|
||||
|
||||
openFile :: FilePath -> IOMode -> IO Handle
|
||||
openFile f mode = do
|
||||
|
||||
@@ -96,7 +96,7 @@ instance MsgStoreClass STMMsgStore where
|
||||
pure msgs
|
||||
|
||||
writeMsg :: STMMsgStore -> STMMsgQueue -> Bool -> Message -> ExceptT ErrorType IO (Maybe (Message, Bool))
|
||||
writeMsg _ STMMsgQueue {msgQueue = q, quota, canWrite, size} _logState !msg = liftIO $ atomically $ do
|
||||
writeMsg _ STMMsgQueue {msgQueue = q, quota, canWrite, size} _logState msg = liftIO $ atomically $ do
|
||||
canWrt <- readTVar canWrite
|
||||
empty <- isEmptyTQueue q
|
||||
if canWrt || empty
|
||||
@@ -105,11 +105,11 @@ instance MsgStoreClass STMMsgStore where
|
||||
writeTVar canWrite $! canWrt'
|
||||
modifyTVar' size (+ 1)
|
||||
if canWrt'
|
||||
then writeTQueue q msg $> Just (msg, empty)
|
||||
else writeTQueue q msgQuota $> Nothing
|
||||
then (writeTQueue q $! msg) $> Just (msg, empty)
|
||||
else (writeTQueue q $! msgQuota) $> Nothing
|
||||
else pure Nothing
|
||||
where
|
||||
!msgQuota = MessageQuota {msgId = msgId msg, msgTs = msgTs msg}
|
||||
msgQuota = MessageQuota {msgId = msgId msg, msgTs = msgTs msg}
|
||||
|
||||
setOverQuota_ :: STMMsgQueue -> IO ()
|
||||
setOverQuota_ q = atomically $ writeTVar (canWrite q) False
|
||||
|
||||
Reference in New Issue
Block a user