[mythtv-users] Storage Groups, free space, (and auto-expire?)

Chris Pinkham cpinkham at bc2va.org
Thu Oct 25 03:10:50 UTC 2007


Replying to mutiple emails in this thread...

I have considered adding code to the storage scheduling logic to take into
account whether a recording will have to be expired by putting a new 
recording onto the same directory, but I rarely have anything expire for
disk space reasons so I haven't looked at this yet.  With the new
'Deleted' recording group functionality, this becomes more important, but
I don't use that feature either.  If all filesystems are full and
something will end up expiring or deleting when create a new recording, 
then we should choose the filesystem with the first expirable/deletable 
recording rather than chosing a filesystem based on the lowest weight.

I think it is best to put the move functionality in the JobQueue and have
a partially complete patch to add a 'Migration' job for just that.  I've
attached the patch I was working on.  It can perform a bandwidth-throttled
move between directories.  There isn't any GUI around this, and I haven't
spent any time on it recently.  If the JobQueue is the one moving the files
around, then it makes other functionality easier to add such as moving files
instead of deleting them in the AutoExpirer.  It also leads to another idea
on my TODO of adding a nightly housekeeper process to schedule jobs to
move files around in order to try to balance utilization across storage 
directories.  'Archiving' files to larger/slower/etc. storage is also
easy and can be automated via a GUI.  If done using a builtin migration 
job, the perl bindings become easy. :)  The attached patch doesn't handle
moving files between servers, all drives have to be accessible locally
(ie, physically local or network mounted).  It does handle moving
files between Storage Groups and does update the storagegroup field in
the DB.

Once the move functionality is there, then I have another UserJob idea 
as well.  I'd like to have the ability to record files to local storage
then have a job that moves the recordings to a larger NAS-mounted
directory later.  This way you can have backend systems with relatively
small local drives for recording onto, but long-term (ie, longer than
24 hours or so) recordings get moved onto larger shared filesystems.
With current Storage Groups functionality, as long as the backend has
the local drive and NFS drive(s) mounted, then the storage scheduler 
would prefer putting new recordings onto the local storage unless it
filled up and then recordings would be written to the NAS directory
(unless you have more than 2 simultaneous recordings when the first
would go to local and the 3rd to NAS with the current weights in the
SG code).

--
Chris
-------------- next part --------------
Index: libs/libmythtv/jobqueue.cpp
===================================================================
--- libs/libmythtv/jobqueue.cpp	(revision 12478)
+++ libs/libmythtv/jobqueue.cpp	(working copy)
@@ -10,6 +10,7 @@
 #include <sys/stat.h>
 #include <sys/wait.h>
 #include <sys/resource.h>
+#include <fcntl.h>
 
 #include <iostream>
 using namespace std;
@@ -23,6 +24,14 @@
 #include "mythdbcon.h"
 #include "previewgenerator.h"
 
+#ifndef O_STREAMING
+#define O_STREAMING 0
+#endif
+
+#ifndef O_LARGEFILE
+#define O_LARGEFILE 0
+#endif
+
 #define LOC     QString("JobQueue: ")
 #define LOC_ERR QString("JobQueue Error: ")
 
@@ -1066,6 +1075,7 @@
     {
         case JOB_TRANSCODE:  return tr("Transcode");
         case JOB_COMMFLAG:   return tr("Flag Commercials");
+        case JOB_MIGRATE:    return tr("Recording File Migration");
     }
 
     if (jobType & JOB_USERJOB)
@@ -1275,6 +1285,7 @@
                                  break;
             case JOB_COMMFLAG:   allowSetting = "JobAllowCommFlag";
                                  break;
+            case JOB_MIGRATE:    return true;
             default:             return false;
         }
     }
@@ -1537,6 +1548,10 @@
     {
         StartChildJob(FlagCommercialsThread, pginfo);
     }
+    else if (jobType == JOB_MIGRATE)
+    {
+        StartChildJob(MigrateThread, pginfo);
+    }
     else if (jobType & JOB_USERJOB)
     {
         StartChildJob(UserJobThread, pginfo);
@@ -1579,6 +1594,8 @@
         return "Transcode";
     else if (jobType == JOB_COMMFLAG)
         return "Commercial Flagging";
+    else if (jobType == JOB_MIGRATE)
+        return "Recording File Migration";
     else if (!(jobType & JOB_USERJOB))
         return "Unknown Job";
 
@@ -1611,6 +1628,10 @@
         if (command == "mythcommflag")
             return command;
     }
+    else if (jobType == JOB_MIGRATE)
+    {
+        return "INTERNAL";
+    }
     else if (jobType & JOB_USERJOB)
     {
         command = gContext->GetSetting(
@@ -1906,12 +1927,7 @@
         ChangeJobStatus(jobID, JOB_ERRORED, "Retry limit exceeded");
     }
 
-    controlFlagsLock.lock();
-    runningJobIDs.erase(key);
-    runningJobTypes.erase(key);
-    runningJobDescs.erase(key);
-    runningJobCommands.erase(key);
-    controlFlagsLock.unlock();
+    ClearRunningJobInfo(key);
 }
 
 void *JobQueue::FlagCommercialsThread(void *param)
@@ -2062,16 +2078,12 @@
         program_info->pathname = program_info->GetPlaybackURL();
         (new PreviewGenerator(program_info, true))->Run();
     }
+    controlFlagsLock.unlock();
 
     if (msg != "")
         VERBOSE(VB_IMPORTANT, LOC + msg);
 
-    jobControlFlags.erase(key);
-    runningJobIDs.erase(key);
-    runningJobTypes.erase(key);
-    runningJobDescs.erase(key);
-    runningJobCommands.erase(key);
-    controlFlagsLock.unlock();
+    ClearRunningJobInfo(key);
 
     delete program_info;
 }
@@ -2158,7 +2170,206 @@
         gContext->dispatch(me);
     }
 
+    ClearRunningJobInfo(key);
+}
+
+void *JobQueue::MigrateThread(void *param)
+{
+    JobQueue *theMigrater = (JobQueue *)param;
+    theMigrater->DoMigrateThread();
+
+    return NULL;
+}
+
+void JobQueue::DoMigrateThread(void)
+{
+    if (!m_pginfo)
+        return;
+
+    ProgramInfo *program_info = new ProgramInfo(*m_pginfo);
+    int controlMigrating = JOB_RUN;
+    QString subtitle = program_info->subtitle.isEmpty() ? "" :
+                           QString(" \"%1\"").arg(program_info->subtitle);
+    QString logDesc = QString("%1%2 recorded from channel %3 at %4")
+                          .arg(program_info->title.local8Bit())
+                          .arg(subtitle.local8Bit())
+                          .arg(program_info->chanid)
+                          .arg(program_info->recstartts.toString());
+    
+    QString key = GetJobQueueKey(program_info);
+    int jobID = runningJobIDs[key];
+
+    childThreadStarted = true;
+
+    QString args = GetJobArgs(jobID);
+    QStringList tokens = QStringList::split(":", args, true);
+    QString storageGroup = "";
+    QString storageDir = "";
+    QString storageFile = "";
+
+    if (tokens.size() >= 3)
+    {
+        storageGroup = tokens[0];
+        storageDir   = tokens[1];
+        storageFile  = tokens[2];
+    }
+    else
+    {
+        ChangeJobStatus(jobID, JOB_ERRORED, "Job Arguments were not supplied.");
+        return;
+    }
+
+    QString src = program_info->GetPlaybackURL(false, true);
+
+    if (src.left(1) != "/")
+    {
+        ChangeJobStatus(jobID, JOB_ERRORED,
+                        "Unable to find file anywhere local.");
+        return;
+    }
+
+    QString dest = storageDir + "/" + storageFile;
+    QString tmpdest = dest + ".tmp";
+
+    ChangeJobStatus(jobID, JOB_RUNNING);
+
     controlFlagsLock.lock();
+    jobControlFlags[key] = &controlMigrating;
+    controlFlagsLock.unlock();
+
+    QString msg = "Recording Migration Starting";
+    VERBOSE(VB_GENERAL, LOC + QString("%1 for %2").arg(msg).arg(logDesc));
+    gContext->LogEntry("migrate", LP_NOTICE, msg, logDesc);
+
+    program_info->MarkAsInUse(true, "filemigration");
+
+    const int readSize = 600000;
+    int ifd, ofd;
+    char buf[readSize];
+    int r;
+    bool ok = true;
+    int sleepTime = 250000;
+    size_t offset = 0;
+    int loop = 0;
+    size_t srcSize = 0;
+
+    struct stat st;
+    if (stat(src.ascii(), &st) == 0)
+        srcSize = st.st_size;
+
+    ifd = open(src.ascii(), O_RDONLY|O_LARGEFILE|O_STREAMING);
+    if (ifd >= 0)
+    {
+        ofd = open(tmpdest.ascii(), O_WRONLY|O_TRUNC|O_CREAT|O_LARGEFILE, 0644);
+        if (ofd < 0)
+        {
+            close(ifd);
+            ok = false;
+            ChangeJobStatus(jobID, JOB_ERRORED, QString(
+                "ERROR: Unable to open temporary destination file '%1'.")
+                .arg(tmpdest));
+        }
+
+        if (ok)
+        {
+            bool done = false;
+            while(ok && !done && ((r = read(ifd, buf, readSize)) > 0))
+            {
+                if (write(ofd, buf, r) != r)
+                {
+                    ok = false;
+                    ChangeJobStatus(jobID, JOB_ERRORED, QString(
+                        "ERROR: Error writing to destination at offset %1.")
+                        .arg(offset));
+                }
+                else
+                {
+                    usleep(sleepTime);
+                }
+                offset += r;
+                loop++;
+
+                if ((loop % 100) == 0)
+                {
+                    if (srcSize)
+                        ChangeJobComment(jobID, QString(
+                            "Migrated %1KB of %2KB (%3%%)")
+                            .arg(offset / 1024).arg(srcSize / 1024)
+                            .arg((offset/1024.0)/(srcSize/1024.0), 4, 'f', 1));
+                    else
+                        ChangeJobComment(jobID, QString(
+                            "Migrated %1KB").arg(offset/1024));
+
+                    program_info->UpdateInUseMark(true);
+                }
+            }
+
+            close(ifd);
+            close(ofd);
+
+            if (srcSize && (offset != srcSize))
+            {
+                ok = false;
+                ChangeJobStatus(jobID, JOB_ERRORED, QString(
+                    "ERROR: Only copied %1 of %2 bytes in file.")
+                    .arg(offset).arg(srcSize));
+            }
+            else
+            {
+                ChangeJobComment(jobID,
+                    "Copied complete file.");
+            }
+        }
+    }
+    else
+    {
+        ChangeJobStatus(jobID, JOB_ERRORED, QString(
+            "ERROR: Unable to open source file '%1'.").arg(src));
+    }
+
+    if (ok && (storageGroup != program_info->storagegroup))
+    {
+        MSqlQuery query(MSqlQuery::InitCon());
+
+        query.prepare("UPDATE recorded SET storagegroup = :STORAGEGROUP "
+                      "WHERE chanid = :CHANID AND starttime = :STARTTIME;");
+        query.bindValue(":CHANID", program_info->chanid);
+        query.bindValue(":STARTTIME", program_info->recstartts);
+        query.bindValue(":STORAGEGROUP", storageGroup);
+
+        query.exec();
+    }
+
+    if (ok)
+    {
+        if (rename(tmpdest.ascii(), dest.ascii()) == 0)
+        {
+            unlink(src.ascii());
+
+            ChangeJobStatus(jobID, JOB_FINISHED, "Migration Complete.");
+        }
+        else
+        {
+            unlink(tmpdest.ascii());
+            ChangeJobStatus(jobID, JOB_ERRORED, 
+                "ERROR: Unable to rename tmp file.");
+        }
+    }
+    else
+    {
+        unlink(tmpdest.ascii());
+    }
+
+    program_info->MarkAsInUse(false);
+    delete program_info;
+
+    ClearRunningJobInfo(key);
+}
+
+void JobQueue::ClearRunningJobInfo(QString key)
+{
+    controlFlagsLock.lock();
+    jobControlFlags.erase(key);
     runningJobIDs.erase(key);
     runningJobTypes.erase(key);
     runningJobDescs.erase(key);
Index: libs/libmythtv/jobqueue.h
===================================================================
--- libs/libmythtv/jobqueue.h	(revision 12478)
+++ libs/libmythtv/jobqueue.h	(working copy)
@@ -69,6 +69,7 @@
     JOB_SYSTEMJOB    = 0x00ff,
     JOB_TRANSCODE    = 0x0001,
     JOB_COMMFLAG     = 0x0002,
+    JOB_MIGRATE      = 0x0004,
 
     JOB_USERJOB      = 0xff00,
     JOB_USERJOB1     = 0x0100,
@@ -185,6 +186,7 @@
 
     QString GetJobDescription(int jobType);
     QString GetJobCommand(int id, int jobType, ProgramInfo *tmpInfo);
+    void ClearRunningJobInfo(QString key);
 
     static void *TranscodeThread(void *param);
     static QString PrettyPrint(off_t bytes);
@@ -196,6 +198,9 @@
     static void *UserJobThread(void *param);
     void DoUserJobThread(void);
 
+    static void *MigrateThread(void *param);
+    void DoMigrateThread(void);
+
     QString m_hostname;
 
     int jobsRunning;


More information about the mythtv-users mailing list