#ifndef LIB9PSTREAMER_H #define LIB9PSTREAMER_H #include #include #include #include #include #include #include "lib9pclient.h" /** * Lib9pStreamer - Streaming file reader using native lib9p. * * Runs blocking reads in a worker thread, emitting data via signals. * Replaces NativeStreamer (which used ollie-9p subprocess). * * Restart policies: * Oneshot — runs once, emits finished on EOF/error * Looping — always restarts after EOF/error * Guarded — restarts only if guard() returns true */ class Lib9pStreamer : public QObject { Q_OBJECT public: enum RestartPolicy { Oneshot, Looping, Guarded }; explicit Lib9pStreamer(RestartPolicy policy, QObject *parent = nullptr); ~Lib9pStreamer() override; void start(const QString &path); void stop(); bool isRunning() const; void setGuard(std::function guard); signals: void dataReady(const QByteArray &data); void finished(); // Only emitted when not restarting void errorOccurred(const QString &error); private slots: void onWorkerData(const QByteArray &data); void onWorkerFinished(bool success, const QString &error); void scheduleRestart(); private: void startWorker(); void stopWorker(); QString m_path; RestartPolicy m_policy; std::function m_guard; QTimer *m_restartTimer = nullptr; bool m_stopping = false; // Worker thread state class Worker; Worker *m_worker = nullptr; QThread *m_workerThread = nullptr; }; // Worker that runs in a separate thread class Lib9pStreamer::Worker : public QObject { Q_OBJECT public: explicit Worker(const QString &path); ~Worker() override; void requestStop(); public slots: void run(); signals: void dataReady(const QByteArray &data); void finished(bool success, const QString &error); private: QString m_path; volatile bool m_stopRequested = false; Lib9pClient *m_client = nullptr; int m_fid = -1; QMutex m_fidMutex; // protects m_fid during close }; #endif // LIB9PSTREAMER_H