Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/c-cpp.yml
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ on:
jobs:
build:

runs-on: ubuntu-latest
runs-on: ubuntu-26.04

steps:
- uses: actions/checkout@v7
Expand Down
29 changes: 29 additions & 0 deletions common/coroutines.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
/* coroutines.h
*
* Copyright (c) 2026 Andrey Kutejko <andy128k@gmail.com>
*
* This program is free software; you can redistribute it and/or
* modify it under the terms of the GNU General Public License as
* published by the Free Software Foundation; either version 2 of the
* License, or (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License
* along with this program; if not, write to the Free Software
* Foundation, Inc., 59 Temple Place, Suite 330, Boston, MA 02111-1307
* USA
*/

#ifndef _COROUTINES_H_
#define _COROUTINES_H_

#include <coroutine>
#include <beman/task/task.hpp>

template <class T = void> using task = beman::execution::task<T>;

#endif
3 changes: 2 additions & 1 deletion common/meson.build
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
common_sources = files(
'racquireasync.cc',
'rconfiguration.cc',
'rinstallprogress.cc',
'rpackage.cc',
Expand All @@ -23,7 +24,7 @@ libsynaptic = static_library(
'synaptic',
common_sources,
cpp_args: rpm_compile_args,
dependencies: common_deps,
dependencies: common_deps + [task_dep],
include_directories: [root_inc, common_inc],
)

Expand Down
177 changes: 177 additions & 0 deletions common/racquireasync.cc
Original file line number Diff line number Diff line change
@@ -0,0 +1,177 @@
#include "config.h" // IWYU pragma: associated

#include "racquireasync.h"

#include <queue>

// Fetcher message

struct MessageMediaChange
{
std::string media;
std::string drive;
};
struct MessageIMSHit
{
pkgAcquire::ItemDesc item;
};
struct MessageFetch
{
pkgAcquire::ItemDesc item;
};
struct MessageDone
{
pkgAcquire::ItemDesc item;
};
struct MessageFail
{
pkgAcquire::ItemDesc item;
};
struct MessageStart
{};
struct MessageStop
{};
struct MessageFinish
{
std::any result;
};
using FetcherMessage = std::variant<MessageMediaChange,
MessageIMSHit,
MessageFetch,
MessageDone,
MessageFail,
MessageStart,
MessageStop,
MessageFinish>;

template <typename T> class EventQueue
{
public:
void push(T event)
{
std::lock_guard<std::mutex> lock(mutex_);
queue_.push(std::move(event));
}

std::optional<T> poll()
{
std::lock_guard<std::mutex> lock(mutex_);

if (queue_.empty())
return std::nullopt;

auto event = std::move(queue_.front());
queue_.pop();
return event;
}

T pop()
{
while (true) {
if (auto answer = poll()) {
return answer.value();
} else {
usleep(100000);
}
}
}

private:
std::mutex mutex_;
std::queue<T> queue_;
};

class pkgAcquireStatusQueue : public pkgAcquireStatus
{
public:
virtual bool MediaChange(std::string Media, std::string Drive) override
{
messageQueue.push(MessageMediaChange{Media, Drive});
return answerQueue.pop();
}
virtual void IMSHit(pkgAcquire::ItemDesc &Itm) override
{
messageQueue.push(MessageIMSHit{Itm});
}
virtual void Fetch(pkgAcquire::ItemDesc &Itm) override
{
messageQueue.push(MessageFetch{Itm});
}
virtual void Done(pkgAcquire::ItemDesc &Itm) override
{
messageQueue.push(MessageDone{Itm});
}
virtual void Fail(pkgAcquire::ItemDesc &Itm) override
{
messageQueue.push(MessageFail{Itm});
}
virtual void Start() override
{
messageQueue.push(MessageStart{});
}
virtual void Stop() override
{
messageQueue.push(MessageStop{});
}

explicit pkgAcquireStatusQueue(EventQueue<FetcherMessage> &messageQueue,
EventQueue<bool> &answerQueue)
: messageQueue(messageQueue), answerQueue(answerQueue)
{}

private:
EventQueue<FetcherMessage> &messageQueue;
EventQueue<bool> &answerQueue;
};

task<std::any> runWithStatusAsyncAny(
std::function<std::any(pkgAcquireStatus &)> &&body,
RPkgAcquireStatusAsync *status)
{
EventQueue<FetcherMessage> messageQueue;
EventQueue<bool> answerQueue;
pkgAcquireStatusQueue statusQueue{messageQueue, answerQueue};

std::thread worker([&]() {
auto result = body(statusQueue);
messageQueue.push(FetcherMessage{MessageFinish{result}});
});

while (true) {
if (auto event = messageQueue.poll()) {
std::optional<std::any> result = co_await std::visit(
[status,
&answerQueue](auto &&event) -> task<std::optional<std::any>> {
using T = std::decay_t<decltype(event)>;
if constexpr (std::is_same_v<T, MessageFinish>) {
co_return std::optional<std::any>{event.result};
}
if constexpr (std::is_same_v<T, MessageMediaChange>) {
bool result =
co_await status->MediaChange(event.media, event.drive);
answerQueue.push(result);
} else if constexpr (std::is_same_v<T, MessageIMSHit>) {
co_await status->IMSHit(event.item);
} else if constexpr (std::is_same_v<T, MessageFetch>) {
co_await status->Fetch(event.item);
} else if constexpr (std::is_same_v<T, MessageDone>) {
co_await status->Done(event.item);
} else if constexpr (std::is_same_v<T, MessageFail>) {
co_await status->Fail(event.item);
} else if constexpr (std::is_same_v<T, MessageStart>) {
co_await status->Start();
} else if constexpr (std::is_same_v<T, MessageStop>) {
co_await status->Stop();
}
co_return std::optional<std::any>{};
},
event.value());
if (result) {
worker.join();
co_return result.value();
}
} else {
// co_await sleep_ms{100};
}
}
}
57 changes: 57 additions & 0 deletions common/racquireasync.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
#pragma once

#include "coroutines.h"
#include <apt-pkg/acquire.h>
#include <any>
#include <functional>

class RPkgAcquireStatusAsync
{
public:
[[nodiscard]] virtual task<bool> MediaChange(std::string Media,
std::string Drive) = 0;
[[nodiscard]] virtual task<void> IMSHit(pkgAcquire::ItemDesc &Itm) = 0;
[[nodiscard]] virtual task<void> Fetch(pkgAcquire::ItemDesc &Itm) = 0;
[[nodiscard]] virtual task<void> Done(pkgAcquire::ItemDesc &Itm) = 0;
[[nodiscard]] virtual task<void> Fail(pkgAcquire::ItemDesc &Itm) = 0;
[[nodiscard]] virtual task<void> Start() = 0;
[[nodiscard]] virtual task<void> Stop() = 0;
};

[[nodiscard]] task<std::any> runWithStatusAsyncAny(
std::function<std::any(pkgAcquireStatus &)> &&body,
RPkgAcquireStatusAsync *status);

template <typename R>
[[nodiscard]] task<R> runWithStatusAsync(
std::function<R(pkgAcquireStatus &)> &&body,
RPkgAcquireStatusAsync *status)
{
auto anyBody = [&](pkgAcquireStatus &status) -> std::any {
return body(status);
};
auto result = co_await runWithStatusAsyncAny(std::move(anyBody), status);
co_return std::any_cast<R>(result);
}

[[nodiscard]] inline task<pkgAcquire::RunResult> acquireRunAsync(
pkgAcquire *acquire,
RPkgAcquireStatusAsync *status,
int PulseInterval = 500000)
{
pkgAcquire::RunResult result =
co_await runWithStatusAsync<pkgAcquire::RunResult>(
[acquire,
PulseInterval](pkgAcquireStatus &status) -> pkgAcquire::RunResult {
acquire->SetLog(&status);
pkgAcquire::RunResult result = acquire->Run(
#ifndef HAVE_RPM
PulseInterval
#endif
);
acquire->SetLog(nullptr);
return result;
},
status);
co_return result;
}
2 changes: 1 addition & 1 deletion common/rinstallprogress.cc
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,7 @@ std::optional<pkgPackageManager::OrderResult> RInstallProgress::start(

res = pm->DoInstallPreFork();
if (res == pkgPackageManager::Failed)
return res;
co_return res;

/*
* This will make a pipe from where we can read child's output
Expand Down
28 changes: 20 additions & 8 deletions common/rinstallprogress.h
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,8 @@
#include <optional>
#include <string>

#include "coroutines.h"

class RInstallProgress
{
protected:
Expand Down Expand Up @@ -64,17 +66,27 @@ class RInstallProgress
int numPackages = 0,
int numPackagesTotal = 0);
virtual std::optional<pkgPackageManager::OrderResult> poll();
virtual void finish()
{}
[[nodiscard]] virtual task<void> finish()
{
co_return;
}

virtual void startUpdate()
{}
virtual void updateInterface()
{}
virtual void finishUpdate()
{}
[[nodiscard]] virtual task<void> startUpdate()
{
co_return;
}
[[nodiscard]] virtual task<void> updateInterface()
{
co_return;
}
[[nodiscard]] virtual task<void> finishUpdate()
{
co_return;
}

RInstallProgress()
: _donePackagesTotal(0), _numPackagesTotal(0), _updateFinished(false)
{}
virtual ~RInstallProgress()
{}
};
15 changes: 9 additions & 6 deletions common/rpackage.cc
Original file line number Diff line number Diff line change
Expand Up @@ -960,7 +960,9 @@ void RPackage::setRemove(bool purge)
_lister->notifyChange(this);
}

string RPackage::getScreenshotFile(pkgAcquire *fetcher, bool thumb)
task<string> RPackage::getScreenshotFile(pkgAcquire *fetcher,
RPkgAcquireStatusAsync *status,
bool thumb)
{
string descr("Screenshot for ");
descr += name();
Expand All @@ -985,12 +987,13 @@ string RPackage::getScreenshotFile(pkgAcquire *fetcher, bool thumb)
new pkgAcqFile(
fetcher, uri, HashStringList(), 0, descr, name(), "", filename);

fetcher->Run();
co_await acquireRunAsync(fetcher, status);

return filename;
co_return filename;
}

string RPackage::getChangelogFile(pkgAcquire *fetcher)
task<string> RPackage::getChangelogFile(pkgAcquire *fetcher,
RPkgAcquireStatusAsync *status)
{
string descr("Changelog for ");
descr += name();
Expand All @@ -1009,7 +1012,7 @@ string RPackage::getChangelogFile(pkgAcquire *fetcher)


ofstream out(filename.c_str());
if (fetcher->Run() == pkgAcquire::Failed) {
if (co_await acquireRunAsync(fetcher, status) == pkgAcquire::Failed) {
out << "Failed to download the list of changes. " << endl;
out << "Please check your Internet connection." << endl;
// FIXME: Need to dequeue the item
Expand All @@ -1029,7 +1032,7 @@ string RPackage::getChangelogFile(pkgAcquire *fetcher)
};
out.close();

return filename;
co_return filename;
}

string RPackage::getCandidateOriginStr()
Expand Down
Loading