Pithos 0.3.22: s3d is stopped gently and never mid-upload; stuck waiting files can be repaired
This commit is contained in:
parent
c2bc27b0a1
commit
9d46e47eb1
7 changed files with 210 additions and 11 deletions
|
|
@ -1,7 +1,7 @@
|
|||
{
|
||||
"id": "pithos",
|
||||
"name": "Pithos",
|
||||
"version": "0.3.21",
|
||||
"version": "0.3.22",
|
||||
"description": "Run s3d, the Sia S3 gateway, from the Theseus sidebar: connect it to a Sia indexer, create S3 users and access keys, browse and share buckets, and watch uploads reach Sia.",
|
||||
"author": "Silent Mode",
|
||||
"icon": "data:image/svg+xml;base64,PHN2ZyB4bWxucz0iaHR0cDovL3d3dy53My5vcmcvMjAwMC9zdmciIHZpZXdCb3g9IjAgMCA2NCA2NCI+PGRlZnM+PGxpbmVhckdyYWRpZW50IGlkPSJnIiB4MT0iMCIgeTE9IjAiIHgyPSIwIiB5Mj0iMSI+PHN0b3Agb2Zmc2V0PSIwIiBzdG9wLWNvbG9yPSIjM2RkYzk3Ii8+PHN0b3Agb2Zmc2V0PSIxIiBzdG9wLWNvbG9yPSIjMWU4ZjZhIi8+PC9saW5lYXJHcmFkaWVudD48L2RlZnM+PHBhdGggZD0iTTIyIDZoMjB2NWMwIDItMiAzLTIgNSA4IDMgMTQgMTEgMTQgMjEgMCAxMi05IDIxLTIyIDIxUzEwIDQ5IDEwIDM3YzAtMTAgNi0xOCAxNC0yMSAwLTItMi0zLTItNXoiIGZpbGw9InVybCgjZykiLz48cGF0aCBkPSJNMTQgMzBoMzZNMTMgNDBoMzgiIHN0cm9rZT0iIzBiMWYxOCIgc3Ryb2tlLXdpZHRoPSIzIiBzdHJva2UtbGluZWNhcD0icm91bmQiIG9wYWNpdHk9Ii41NSIvPjwvc3ZnPgo=",
|
||||
|
|
|
|||
|
|
@ -4,6 +4,7 @@
|
|||
import { spawn, execFile } from 'node:child_process';
|
||||
import { EventEmitter } from 'node:events';
|
||||
import path from 'node:path';
|
||||
import { sendBreak } from './winctrl.js';
|
||||
|
||||
const LOG_LINES = 2000;
|
||||
const ANSI = /\x1b\[[0-9;]*m/g;
|
||||
|
|
@ -19,6 +20,8 @@ export class Daemon extends EventEmitter {
|
|||
this.startedAt = null;
|
||||
this.logs = [];
|
||||
this.seq = 0;
|
||||
this.busy = 0; // flushes Pithos started and is waiting on
|
||||
this.lastUploadAt = 0; // s3d's last "uploading object group"
|
||||
}
|
||||
|
||||
env() {
|
||||
|
|
@ -54,6 +57,7 @@ export class Daemon extends EventEmitter {
|
|||
// s3d logs "server started" once every listener is up.
|
||||
if (this.state === 'starting' && /server started/.test(line)) this.setState('running');
|
||||
if (/No app key found/.test(line)) this.needsLogin = true;
|
||||
if (/uploading object group/.test(line)) this.lastUploadAt = Date.now();
|
||||
this.logs.push(entry);
|
||||
if (this.logs.length > LOG_LINES) this.logs.shift();
|
||||
this.emit('log', entry);
|
||||
|
|
@ -101,14 +105,49 @@ export class Daemon extends EventEmitter {
|
|||
return this.status();
|
||||
}
|
||||
|
||||
async stop() {
|
||||
// Work s3d must not be stopped in the middle of (Upload now, autoflush).
|
||||
async whileBusy(fn) {
|
||||
this.busy++;
|
||||
try { return await fn(); } finally { this.busy--; }
|
||||
}
|
||||
|
||||
// True once no flush of ours is running and s3d has not started an upload
|
||||
// batch for a while (it logs the start of one, not the end).
|
||||
quiet() { return this.busy === 0 && Date.now() - this.lastUploadAt > 120_000; }
|
||||
|
||||
async waitForQuiet(maxMs) {
|
||||
const until = Date.now() + maxMs;
|
||||
let told = false;
|
||||
while (!this.quiet() && Date.now() < until && this.child) {
|
||||
if (!told) { this.pushLog('sys', 'waiting for the upload in progress to finish before stopping s3d'); told = true; }
|
||||
await new Promise((r) => setTimeout(r, 1000));
|
||||
}
|
||||
return this.quiet();
|
||||
}
|
||||
|
||||
// An s3d stopped in the middle of an upload has lost buffered files whose
|
||||
// rows still said "waiting", so a stop first lets uploads finish, then
|
||||
// asks s3d to shut down cleanly (Ctrl+Break on Windows, SIGTERM elsewhere)
|
||||
// and only ends it by force if it is still there a minute later.
|
||||
// waitMs: how long to wait for uploads (short when the host is quitting).
|
||||
async stop({ waitMs = 5 * 60_000 } = {}) {
|
||||
const child = this.child;
|
||||
if (!child && this.pendingStart) { this.pendingStart = false; this.setState('stopped'); return this.status(); }
|
||||
if (!child) return this.status();
|
||||
this.setState('stopping');
|
||||
const exited = new Promise((r) => child.once('exit', r));
|
||||
child.kill(); // SIGTERM on unix; TerminateProcess on Windows
|
||||
const timer = setTimeout(() => { try { child.kill('SIGKILL'); } catch {} }, 10_000);
|
||||
const exited = new Promise((r) => { if (this.child !== child) r(); else child.once('exit', r); });
|
||||
if (!(await this.waitForQuiet(waitMs)) && this.child === child) this.pushLog('sys', 'an upload is still running; stopping s3d anyway');
|
||||
if (this.child === child) {
|
||||
let asked = false;
|
||||
if (process.platform === 'win32') asked = await sendBreak(child.pid);
|
||||
else { try { asked = child.kill('SIGTERM'); } catch {} }
|
||||
if (!asked) { try { child.kill(); } catch {} }
|
||||
}
|
||||
const timer = setTimeout(() => {
|
||||
if (this.child !== child) return;
|
||||
this.pushLog('sys', 's3d did not stop within a minute; ending it');
|
||||
try { child.kill('SIGKILL'); } catch {}
|
||||
}, 60_000);
|
||||
await exited;
|
||||
clearTimeout(timer);
|
||||
return this.status();
|
||||
|
|
|
|||
|
|
@ -29,6 +29,7 @@ import { accountInfo, readConnection, fingerprint } from './sia-account.js';
|
|||
import { guessType } from './hosted.js';
|
||||
import { accountSpace } from './space.js';
|
||||
import { createShareUrl } from './sia-share.js';
|
||||
import { stuckUploads } from './stuck.js';
|
||||
|
||||
const HERE = path.dirname(fileURLToPath(import.meta.url));
|
||||
const UI_DIR = path.join(HERE, '..', 'ui');
|
||||
|
|
@ -73,6 +74,8 @@ export async function createPithos(opts = {}) {
|
|||
|
||||
const configFile = resolveConfigPath(configOpt);
|
||||
const daemon = new Daemon({ configFile, binary: resolveBinary(binaryOpt, binDir) });
|
||||
// a flush s3d is in the middle of keeps any stop waiting (see Daemon.stop)
|
||||
const flushNow = () => daemon.whileBusy(() => admin.flush(cfg()));
|
||||
const cli = makeCli(daemon);
|
||||
const login = new LoginSession(daemon, {
|
||||
registration: () => cli.registration(),
|
||||
|
|
@ -366,7 +369,30 @@ export async function createPithos(opts = {}) {
|
|||
});
|
||||
|
||||
route('GET', '/api/stats', () => admin.uploadStats(cfg()));
|
||||
route('POST', '/api/flush', async () => { await admin.flush(cfg()); return { ok: true }; });
|
||||
route('POST', '/api/flush', async () => {
|
||||
try { await flushNow(); } catch (e) {
|
||||
if (/failed to open upload|cannot find the file|no such file/i.test(e.message)) {
|
||||
throw new HttpError(409, 'Some files waiting for Sia lost their copy on this computer, so the upload stopped. Pithos lists them on the Drives page: upload them again or remove them.');
|
||||
}
|
||||
throw e;
|
||||
}
|
||||
return { ok: true };
|
||||
});
|
||||
|
||||
// Waiting files whose buffered copy is gone (see core/stuck.js).
|
||||
route('GET', '/api/uploads/stuck', async () => {
|
||||
if (hosted.isHosted()) return { items: [] };
|
||||
if (!daemon.quiet()) return { items: [], busy: true };
|
||||
return { items: await stuckUploads(cfg().directory) };
|
||||
});
|
||||
route('POST', '/api/uploads/stuck/remove', async (req) => {
|
||||
const { bucket, key } = await jsonBody(req);
|
||||
const item = (await stuckUploads(cfg().directory)).find((i) => i.bucket === bucket && i.key === key);
|
||||
if (!item) throw new HttpError(409, 'This file is no longer stuck (it may have been uploaded again).');
|
||||
if (!item.user) throw new HttpError(409, 'No S3 user owns this drive.');
|
||||
await (await s3For(item.user)).deleteObject(bucket, key);
|
||||
return { ok: true };
|
||||
});
|
||||
|
||||
const autoFlush = createAutoFlush({
|
||||
settings: () => {
|
||||
|
|
@ -378,7 +404,7 @@ export async function createPithos(opts = {}) {
|
|||
// Only an s3d that is up can be asked; a stopped one just means no flush.
|
||||
stats: () => (daemon.state === 'starting' || daemon.state === 'stopping'
|
||||
? Promise.reject(new Error('busy')) : admin.uploadStats(cfg())),
|
||||
flush: () => admin.flush(cfg()),
|
||||
flush: () => flushNow(),
|
||||
log: (line) => daemon.pushLog('sys', line),
|
||||
});
|
||||
route('GET', '/api/autoflush', () => autoFlush.status());
|
||||
|
|
@ -511,7 +537,7 @@ export async function createPithos(opts = {}) {
|
|||
try { return await createShareUrl(dir, bucket, key, until); }
|
||||
catch (e) {
|
||||
if (e.code !== 'pending') throw new HttpError(e.status || 502, e.message);
|
||||
await admin.flush(cfg()); // send what is waiting, then try once more
|
||||
await flushNow(); // send what is waiting, then try once more
|
||||
try { return await createShareUrl(dir, bucket, key, until); }
|
||||
catch (e2) { throw new HttpError(e2.status || 502, e2.code === 'pending' ? 'This file is still uploading to Sia. Try again in a minute.' : e2.message); }
|
||||
}
|
||||
|
|
@ -949,7 +975,7 @@ export async function createPithos(opts = {}) {
|
|||
login.cancel();
|
||||
for (const s of subscribers) s.end();
|
||||
handoff?.close();
|
||||
await daemon.stop();
|
||||
await daemon.stop({ waitMs: 20_000 }); // the host is quitting: a short wait, then a clean stop
|
||||
// Idle keep-alive sockets would hold close() for seconds; nobody is left to answer.
|
||||
const closed = new Promise((r) => server.close(r));
|
||||
server.closeAllConnections?.();
|
||||
|
|
|
|||
41
bundled-addons/pithos/core/stuck.js
Normal file
41
bundled-addons/pithos/core/stuck.js
Normal file
|
|
@ -0,0 +1,41 @@
|
|||
// Files s3d still lists as waiting for Sia whose buffered copy is gone.
|
||||
//
|
||||
// s3d keeps every not-yet-uploaded object as uploads/<filename> and marks it
|
||||
// uploaded by setting sia_object_id. An s3d stopped in the middle of an
|
||||
// upload has been seen to delete the copies without recording the upload;
|
||||
// every flush after that fails on them ("failed to open upload"). They hold
|
||||
// no data any more: not on this computer, never on Sia. Pithos lists them so
|
||||
// the user can upload the file again or remove the entry.
|
||||
|
||||
import fs from 'node:fs';
|
||||
import path from 'node:path';
|
||||
|
||||
async function candidates(dataDir, olderThanMs) {
|
||||
const file = path.join(dataDir, 's3d.db');
|
||||
if (!fs.existsSync(file)) return [];
|
||||
const { DatabaseSync } = await import('node:sqlite');
|
||||
const db = new DatabaseSync(file, { readOnly: true });
|
||||
try {
|
||||
const rows = db.prepare(`SELECT b.name AS bucket, u.name AS user, o.name AS key, o.size AS size, o.filename AS filename, o.updated_at AS updated
|
||||
FROM objects o JOIN buckets b ON b.id = o.bucket_id LEFT JOIN users u ON u.id = b.user_id
|
||||
WHERE o.sia_object_id IS NULL AND o.is_latest = 1 AND o.is_delete_marker = 0
|
||||
AND o.filename IS NOT NULL AND o.filename != ''`).all();
|
||||
const uploads = path.join(dataDir, 'uploads');
|
||||
const cutoff = Date.now() - olderThanMs;
|
||||
return rows
|
||||
// updated_at is Unix seconds
|
||||
.filter((r) => { const n = Number(r.updated); const t = Number.isFinite(n) ? (n < 1e12 ? n * 1000 : n) : Date.parse(r.updated); return !Number.isFinite(t) || t < cutoff; })
|
||||
.filter((r) => !fs.existsSync(path.join(uploads, path.basename(String(r.filename)))))
|
||||
.map((r) => ({ bucket: r.bucket, user: r.user, key: r.key, size: Number(r.size), filename: String(r.filename) }));
|
||||
} finally { db.close(); }
|
||||
}
|
||||
|
||||
// Only entries missing twice, a moment apart: s3d removes a copy just before
|
||||
// it records the upload, and a list taken in between would be wrong.
|
||||
export async function stuckUploads(dataDir, { olderThanMs = 120_000, settleMs = 2000 } = {}) {
|
||||
const first = await candidates(dataDir, olderThanMs);
|
||||
if (!first.length) return [];
|
||||
await new Promise((r) => setTimeout(r, settleMs));
|
||||
const again = new Set((await candidates(dataDir, olderThanMs)).map((r) => `${r.bucket}\n${r.key}`));
|
||||
return first.filter((r) => again.has(`${r.bucket}\n${r.key}`));
|
||||
}
|
||||
33
bundled-addons/pithos/core/winctrl.js
Normal file
33
bundled-addons/pithos/core/winctrl.js
Normal file
|
|
@ -0,0 +1,33 @@
|
|||
// A gentle stop for s3d on Windows. child.kill() there is TerminateProcess:
|
||||
// s3d dies wherever it is, and one stopped in the middle of an upload had
|
||||
// lost buffered files whose database rows still said "waiting". s3d shuts
|
||||
// down cleanly on an interrupt, so this attaches a short PowerShell to s3d's
|
||||
// own (hidden) console and sends Ctrl+Break there.
|
||||
//
|
||||
// Ctrl+Break, not Ctrl+C: Windows disables Ctrl+C for a process started in
|
||||
// a new process group, and Go reads both as os.Interrupt. A spawned console
|
||||
// program gets a console of its own (Node passes CREATE_NO_WINDOW), so the
|
||||
// event reaches s3d and not Pithos. It ends the PowerShell too, which is
|
||||
// why only "could not attach" (2) and "could not send" (3) count as failure.
|
||||
|
||||
import { execFile } from 'node:child_process';
|
||||
|
||||
const SCRIPT = (pid) => [
|
||||
"Add-Type -Name K -Namespace PithosCtrl -MemberDefinition '",
|
||||
'[DllImport("kernel32.dll")] public static extern bool FreeConsole();',
|
||||
'[DllImport("kernel32.dll")] public static extern bool AttachConsole(uint pid);',
|
||||
'[DllImport("kernel32.dll")] public static extern bool GenerateConsoleCtrlEvent(uint ev, uint group);',
|
||||
"';",
|
||||
'[PithosCtrl.K]::FreeConsole() | Out-Null;',
|
||||
`if (-not [PithosCtrl.K]::AttachConsole(${Number(pid) >>> 0})) { exit 2 };`,
|
||||
'if ([PithosCtrl.K]::GenerateConsoleCtrlEvent(1, 0)) { exit 0 } else { exit 3 }',
|
||||
].join(' ');
|
||||
|
||||
// Resolves true when Ctrl+Break was sent to the console of `pid`.
|
||||
export function sendBreak(pid) {
|
||||
if (process.platform !== 'win32' || !pid) return Promise.resolve(false);
|
||||
return new Promise((resolve) => {
|
||||
execFile('powershell.exe', ['-NoLogo', '-NoProfile', '-NonInteractive', '-ExecutionPolicy', 'Bypass', '-Command', SCRIPT(pid)],
|
||||
{ windowsHide: true, timeout: 20_000 }, (err) => resolve(!err || (err.code !== 2 && err.code !== 3 && !err.killed)));
|
||||
});
|
||||
}
|
||||
|
|
@ -475,3 +475,9 @@ a.brand:hover > span:first-of-type { color: var(--accent); }
|
|||
.link-card .row { gap: 6px; flex-wrap: wrap; }
|
||||
.link-make { gap: 10px; align-items: center; flex-wrap: wrap; }
|
||||
.link-card a.btn { text-decoration: none; }
|
||||
|
||||
/* Waiting files that lost their copy (stuck uploads) */
|
||||
.banner.stuck .stack { gap: 8px; width: 100%; }
|
||||
.stuck-list { list-style: none; margin: 0; padding: 0; display: flex; flex-direction: column; gap: 6px; }
|
||||
.stuck-list li { display: flex; gap: 8px; align-items: center; justify-content: space-between; flex-wrap: wrap; }
|
||||
.stuck-name { min-width: 0; overflow: hidden; text-overflow: ellipsis; white-space: nowrap; flex: 1 1 220px; }
|
||||
|
|
|
|||
|
|
@ -294,7 +294,7 @@ function overview(main) {
|
|||
main.append(
|
||||
h('div', { class: 'page-head' }, h('div', {}, h('h1', {}, 'Overview'), h('p', { class: 'muted' }, 'Your s3d gateway and its upload pipeline to Sia.'))),
|
||||
banners,
|
||||
pendingBanner(),
|
||||
stuckBanner(), pendingBanner(),
|
||||
h('div', { class: 'grid two' }, daemonCard, statsCard),
|
||||
);
|
||||
|
||||
|
|
@ -1424,6 +1424,60 @@ function siaAccountBlock(a, onRefresh) {
|
|||
// "N files are waiting to upload to Sia" with Upload now. s3d holds small
|
||||
// uploads until about 40 MB has built up; until then they exist only on this
|
||||
// computer, which is easy to miss.
|
||||
// Files s3d still lists as waiting whose copy on this computer is gone (an
|
||||
// upload was cut off): each can be uploaded again from the original, or removed.
|
||||
function stuckBanner() {
|
||||
const el = h('div', { hidden: true });
|
||||
async function paint() {
|
||||
let items = [];
|
||||
try { items = (await api('GET', '/api/uploads/stuck')).items || []; } catch {}
|
||||
if (!items.length) { el.hidden = true; return; }
|
||||
const remove = async (i) => {
|
||||
try { await api('POST', '/api/uploads/stuck/remove', { bucket: i.bucket, key: i.key }); } catch (e) { fail(e); }
|
||||
};
|
||||
const again = (i) => {
|
||||
const input = h('input', { type: 'file', hidden: true, onchange: async () => {
|
||||
const file = input.files[0];
|
||||
if (!file) return;
|
||||
const path = `/api/s3/${encodeURIComponent(i.user)}/buckets/${encodeURIComponent(i.bucket)}/object?key=${encodeURIComponent(i.key)}`;
|
||||
try {
|
||||
if (bridge) {
|
||||
const r = await bridge.invoke('upload', { path, bytes: await file.arrayBuffer(), type: file.type });
|
||||
if (r.status >= 300) throw new Error(r.data?.error || `HTTP ${r.status}`);
|
||||
} else {
|
||||
const r = await fetch(path, { method: 'PUT', headers: { 'x-pithos': '1', ...(file.type ? { 'content-type': file.type } : {}) }, body: file });
|
||||
if (!r.ok) throw new Error((await r.json().catch(() => ({}))).error || `HTTP ${r.status}`);
|
||||
}
|
||||
toast(`Uploaded ${i.key.split('/').pop()} again`);
|
||||
paint();
|
||||
} catch (e) { fail(e); }
|
||||
} });
|
||||
document.body.append(input);
|
||||
input.click();
|
||||
setTimeout(() => input.remove(), 60_000);
|
||||
};
|
||||
put(el, h('div', { class: 'banner warn stuck' },
|
||||
h('div', { class: 'stack' },
|
||||
h('span', {}, h('strong', {}, `${items.length} file${items.length === 1 ? '' : 's'} lost ${items.length === 1 ? 'its' : 'their'} copy before reaching Sia. `),
|
||||
h('span', { class: 'small' }, items.length === 1
|
||||
? 'An upload was cut off, and it was removed from this computer before it reached Sia. Upload it again from the original, or remove it; until then Upload now stops on it.'
|
||||
: 'An upload was cut off, and they were removed from this computer before they reached Sia. Upload each one again from the original, or remove it; until then Upload now stops on them.')),
|
||||
h('ul', { class: 'stuck-list' }, items.map((i) => h('li', {},
|
||||
h('span', { class: 'stuck-name', title: `${i.bucket}/${i.key}` }, h('span', { class: 'muted' }, `${i.bucket} / `), i.key, h('span', { class: 'muted small' }, ` · ${fmtBytes(i.size)}`)),
|
||||
h('span', { class: 'row' },
|
||||
h('button', { class: 'btn small', onclick: () => again(i) }, 'Upload again…'),
|
||||
h('button', { class: 'btn small danger', onclick: async () => { await remove(i); paint(); } }, 'Remove'))))),
|
||||
items.length > 1 && h('div', {}, h('button', { class: 'btn small danger', onclick: async () => {
|
||||
if (!(await confirmModal(`Remove ${items.length} entries?`, 'Their files are already gone from this computer and never reached Sia; this only removes the names.', 'Remove all'))) return;
|
||||
for (const i of items) await remove(i);
|
||||
paint();
|
||||
} }, 'Remove all')))));
|
||||
el.hidden = false;
|
||||
}
|
||||
paint();
|
||||
return el;
|
||||
}
|
||||
|
||||
function pendingBanner() {
|
||||
const el = h('div', { hidden: true });
|
||||
const st = daemonState();
|
||||
|
|
@ -1494,7 +1548,7 @@ async function bucketList(head, body, user, userSel) {
|
|||
h('div', { class: 'row' }, userPicker(userSel), h('button', { class: 'btn primary', onclick: create }, 'New drive')));
|
||||
const card = h('div', { class: 'card' }, h('p', { class: 'muted' }, 'Loading…'));
|
||||
const bar = h('div', { class: 'view-bar' });
|
||||
put(body, pendingBanner(), bar, card);
|
||||
put(body, stuckBanner(), pendingBanner(), bar, card);
|
||||
const base = `/api/s3/${encodeURIComponent(user)}/buckets`;
|
||||
let mode = viewMode('drives', 'icons');
|
||||
let list = [];
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue